#[cfg(any(feature = "acp", feature = "acp-http"))]
use std::path::PathBuf;
#[cfg(feature = "acp")]
use parking_lot::RwLock;
#[cfg(feature = "acp")]
use crate::agent_setup;
#[cfg(any(feature = "acp", feature = "acp-http"))]
use crate::bootstrap::{AppBuilder, create_mcp_registry};
#[cfg(feature = "acp")]
use zeph_core::agent::Agent;
#[cfg(feature = "acp")]
use zeph_core::channel::Channel;
#[cfg(feature = "acp")]
use zeph_tools::ErasedToolExecutor;
#[cfg(feature = "acp")]
fn resolve_runtime_path(path: &std::path::Path, cwd: &std::path::Path) -> std::path::PathBuf {
if path.is_absolute() {
path.to_path_buf()
} else {
cwd.join(path)
}
}
#[cfg(any(feature = "acp", feature = "acp-http"))]
async fn resolve_acp_auth_clients(
acp_config: &zeph_config::AcpConfig,
vault: &dyn zeph_core::vault::VaultProvider,
) -> anyhow::Result<Vec<zeph_acp::AcpClientToken>> {
let mut clients = Vec::new();
let mut seen_tokens: std::collections::HashSet<String> = std::collections::HashSet::new();
if let Some(ref token) = acp_config.auth_token {
seen_tokens.insert(token.clone());
clients.push(zeph_acp::AcpClientToken {
id: zeph_config::ACP_AUTH_CLIENT_ID_DEFAULT.to_owned(),
token: token.clone(),
});
}
for client in &acp_config.auth_clients {
let token = if let Some(ref t) = client.token {
Some(t.clone())
} else if let Some(ref key) = client.token_vault_key {
match vault.get_secret(key).await {
Ok(Some(t)) if !t.trim().is_empty() => Some(t),
Ok(Some(_)) => {
tracing::warn!(
id = %client.id, vault_key = %key,
"acp.auth_clients: vault key resolved to an empty token; client disabled"
);
None
}
Ok(None) => {
tracing::warn!(
id = %client.id, vault_key = %key,
"acp.auth_clients: vault key not found; client disabled"
);
None
}
Err(e) => {
tracing::warn!(
id = %client.id, vault_key = %key, error = %e,
"acp.auth_clients: failed to resolve token from vault; client disabled"
);
None
}
}
} else {
None
};
let Some(token) = token else { continue };
anyhow::ensure!(
seen_tokens.insert(token.clone()),
"[[acp.auth_clients]] id {:?} resolves to a token that collides with another \
configured client's token",
client.id
);
clients.push(zeph_acp::AcpClientToken {
id: client.id.clone(),
token,
});
}
let auth_declared = acp_config.auth_token.is_some() || !acp_config.auth_clients.is_empty();
anyhow::ensure!(
!auth_declared || !clients.is_empty(),
"[acp] auth_token / [[acp.auth_clients]] configured authentication, but every entry \
failed to resolve (missing vault key, empty vault secret, or vault backend error) — \
refusing to start with authentication silently disabled. Fix the vault key(s), or \
remove the auth_token/auth_clients configuration entirely to run intentionally \
without authentication."
);
Ok(clients)
}
#[cfg(feature = "acp")]
fn log_acp_runtime_paths(config: &zeph_core::config::Config, config_path: &std::path::Path) {
let cwd = std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from("."));
let logging_file = if config.logging.file.is_empty() {
None
} else {
Some(resolve_runtime_path(
std::path::Path::new(&config.logging.file),
&cwd,
))
};
let sqlite_path = resolve_runtime_path(std::path::Path::new(&config.memory.sqlite_path), &cwd);
let debug_output_dir = resolve_runtime_path(config.debug.output_dir.as_path(), &cwd);
let skill_paths: Vec<std::path::PathBuf> = config
.skills
.paths
.iter()
.map(|p| resolve_runtime_path(std::path::Path::new(p), &cwd))
.collect();
let permission_file = config
.acp
.permission_file
.as_ref()
.map(|p| resolve_runtime_path(p.as_path(), &cwd));
tracing::info!(
cwd = %cwd.display(),
config_path = %config_path.display(),
logging_file = logging_file
.as_ref()
.map_or_else(|| "<disabled>".to_owned(), |p| p.display().to_string()),
sqlite_path = %sqlite_path.display(),
debug_output_dir = %debug_output_dir.display(),
permission_file = permission_file
.as_ref()
.map_or_else(|| "<none>".to_owned(), |p| p.display().to_string()),
skill_paths = ?skill_paths,
"ACP startup runtime paths"
);
}
#[cfg(any(feature = "session", feature = "acp"))]
pub(crate) struct SharedCore {
pub(crate) provider: zeph_llm::any::AnyProvider,
pub(crate) embedding_provider: zeph_llm::any::AnyProvider,
pub(crate) registry: std::sync::Arc<parking_lot::RwLock<zeph_skills::registry::SkillRegistry>>,
pub(crate) matcher: Option<zeph_skills::matcher::SkillMatcherBackend>,
pub(crate) memory: std::sync::Arc<zeph_memory::semantic::SemanticMemory>,
pub(crate) budget_tokens: usize,
pub(crate) rl_head: Option<zeph_skills::rl_head::RoutingHead>,
}
#[cfg(any(feature = "session", feature = "acp"))]
pub(crate) async fn build_shared_core(
app: &crate::bootstrap::AppBuilder,
supervisor: &zeph_common::TaskSupervisor,
) -> anyhow::Result<SharedCore> {
let (provider, _status_tx, _status_rx) = app.build_provider().await?;
let embedding_provider = crate::bootstrap::create_embedding_provider(app.config(), &provider);
let budget_tokens = app.auto_budget_tokens(&provider);
let registry = std::sync::Arc::new(parking_lot::RwLock::new(if app.config().cli.safe_mode {
zeph_skills::registry::SkillRegistry::empty()
} else {
app.build_registry()
}));
let memory = std::sync::Arc::new(app.build_memory(&provider, supervisor).await?);
let all_meta_owned: Vec<zeph_skills::loader::SkillMeta> =
registry.read().all_meta().into_iter().cloned().collect();
let all_meta_refs: Vec<&zeph_skills::loader::SkillMeta> = all_meta_owned.iter().collect();
let matcher = app
.build_skill_matcher(&embedding_provider, &all_meta_refs, &memory)
.await;
app.seed_skill_trust_db(&all_meta_owned, &memory).await;
let rl_embed_dim_resolved = if app.config().skills.rl_routing_enabled {
Some(
crate::runner::resolve_rl_embed_dim(
&app.config().skills,
&embedding_provider,
app.config().timeouts.embedding_seconds,
)
.await,
)
} else {
None
};
let rl_head = if let Some(dim) = rl_embed_dim_resolved {
Some(
crate::runner::load_rl_head(&memory)
.await
.unwrap_or_else(|| {
tracing::info!(dim, "rl_head: cold start, initializing fresh routing head");
zeph_skills::rl_head::RoutingHead::new(dim)
}),
)
} else {
None
};
Ok(SharedCore {
provider,
embedding_provider,
registry,
matcher,
memory,
budget_tokens,
rl_head,
})
}
#[cfg(feature = "acp")]
#[allow(clippy::struct_excessive_bools)]
pub(crate) struct SharedAgentDeps {
provider: zeph_llm::any::AnyProvider,
embedding_provider: zeph_llm::any::AnyProvider,
registry: std::sync::Arc<RwLock<zeph_skills::registry::SkillRegistry>>,
matcher: Option<zeph_skills::matcher::SkillMatcherBackend>,
max_active_skills: usize,
skill_disambiguation_threshold: f32,
skill_two_stage_matching: bool,
skill_confusability_threshold: f32,
skill_group_structured: bool,
skill_support_similarity_threshold: f32,
skill_min_injection_score: f32,
skill_generation_provider: String,
skill_disambiguate_provider: String,
semantic_scan: bool,
semantic_scan_provider: String,
trust_config: zeph_core::config::TrustConfig,
rl_routing_enabled: bool,
rl_learning_rate: f32,
rl_weight: f32,
rl_persist_interval: u32,
rl_warmup_updates: u32,
rl_head: Option<zeph_skills::rl_head::RoutingHead>,
tool_executor: std::sync::Arc<dyn zeph_tools::ErasedToolExecutor>,
clock: std::sync::Arc<dyn zeph_common::ClockSource>,
permission_policy: zeph_tools::PermissionPolicy,
policy_gate_pieces: agent_setup::PolicyGatePieces,
capability_scopes_config: zeph_config::CapabilityScopesConfig,
shadow_sentinel_config: zeph_config::ShadowSentinelConfig,
shadow_sentinel_probe_provider: zeph_llm::any::AnyProvider,
trajectory_sentinel_config: zeph_config::TrajectorySentinelConfig,
quality_pipeline: Option<std::sync::Arc<zeph_core::quality::SelfCheckPipeline>>,
skill_paths: Vec<PathBuf>,
pub(crate) memory: std::sync::Arc<zeph_memory::semantic::SemanticMemory>,
history_limit: u32,
recall_limit: usize,
summarization_threshold: usize,
shutdown_summary: bool,
shutdown_summary_min_messages: usize,
shutdown_summary_max_messages: usize,
shutdown_summary_timeout_secs: u64,
shutdown_summary_provider: String,
channel_provider_persistence: bool,
channel_persist_provider_overrides: bool,
index_config: zeph_core::config::IndexConfig,
code_index_provider: zeph_llm::any::AnyProvider,
code_qdrant_ops: Option<zeph_memory::QdrantOps>,
skill_reload_tx: tokio::sync::broadcast::Sender<zeph_skills::watcher::SkillEvent>,
config_reload_tx: tokio::sync::broadcast::Sender<zeph_core::config_watcher::ConfigEvent>,
shutdown_rx: tokio::sync::watch::Receiver<bool>,
config_path: PathBuf,
mcp_tools: Vec<zeph_mcp::McpTool>,
mcp_registry: Option<zeph_mcp::McpToolRegistry>,
mcp_manager: std::sync::Arc<zeph_mcp::McpManager>,
mcp_shared_tools: std::sync::Arc<RwLock<Vec<zeph_mcp::McpTool>>>,
mcp_config: zeph_core::config::McpConfig,
summary_provider: Option<zeph_llm::any::AnyProvider>,
judge_provider: Option<zeph_llm::any::AnyProvider>,
feedback_classifier: Option<zeph_llm::classifier::llm::LlmClassifier>,
#[cfg(feature = "classifiers")]
classifiers_config: zeph_core::config::ClassifiersConfig,
#[cfg(feature = "classifiers")]
pii_filter_enabled: bool,
causal_ipi_config: zeph_sanitizer::causal_ipi::CausalIpiConfig,
causal_provider: Option<zeph_llm::any::AnyProvider>,
nli_config: zeph_sanitizer::nli::NliConfig,
nli_provider: Option<zeph_llm::any::AnyProvider>,
secret_registry: Option<std::sync::Arc<zeph_sanitizer::secret_mask::SecretMaskRegistry>>,
vigil_config: zeph_config::VigilConfig,
probe_provider: Option<zeph_llm::any::AnyProvider>,
planner_provider: Option<zeph_llm::any::AnyProvider>,
verify_provider: Option<zeph_llm::any::AnyProvider>,
ensemble_members: Vec<(String, zeph_llm::any::AnyProvider)>,
orchestrator_provider: Option<zeph_llm::any::AnyProvider>,
predicate_provider: Option<zeph_llm::any::AnyProvider>,
quarantine_provider: Option<(zeph_llm::any::AnyProvider, zeph_sanitizer::QuarantineConfig)>,
guardrail_provider: Option<(
zeph_llm::any::AnyProvider,
zeph_sanitizer::guardrail::GuardrailConfig,
)>,
audit_logger: Option<std::sync::Arc<zeph_tools::AuditLogger>>,
session_config: zeph_core::AgentSessionConfig,
session_persistence_config: zeph_config::SessionConfig,
resume_condenser: zeph_session::LlmCondenser,
resume_token_counter: std::sync::Arc<zeph_agent_context::memory_backend::TokenCounterAdapter>,
provider_pool: Vec<zeph_core::config::ProviderEntry>,
provider_config_snapshot: zeph_core::ProviderConfigSnapshot,
focus_config: zeph_core::config::FocusConfig,
sidequest_config: zeph_core::config::SidequestConfig,
trajectory_config: zeph_core::config::TrajectoryConfig,
category_config: zeph_core::config::CategoryConfig,
tool_filter_config: zeph_core::config::ToolFilterConfig,
hooks_config: zeph_core::config::HooksConfig,
safe_mode: bool,
cwd_allowed_paths: Vec<std::path::PathBuf>,
tools_enabled: bool,
acp_agent_name: String,
acp_agent_version: String,
acp_max_sessions: usize,
acp_session_idle_timeout_secs: u64,
acp_permission_file: Option<std::path::PathBuf>,
acp_available_models: std::sync::Arc<RwLock<Vec<String>>>,
acp_auth_clients: Vec<zeph_acp::AcpClientToken>,
acp_discovery_enabled: bool,
acp_title_max_chars: usize,
acp_max_history: usize,
acp_log_file: Option<String>,
sqlite_path: String,
#[cfg(feature = "acp")]
acp_provider_factory: Option<zeph_acp::ProviderFactory>,
acp_provider_names: Vec<(String, zeph_acp::LlmProtocol)>,
acp_project_rules: Vec<PathBuf>,
acp_additional_directories: Vec<zeph_core::config::AdditionalDir>,
acp_auth_methods: Vec<zeph_core::config::AcpAuthMethod>,
acp_message_ids_enabled: bool,
acp_timeouts: zeph_config::AcpTimeoutsConfig,
acp_model_config: zeph_config::AcpModelConfigConfig,
plugin_dirs_supplier: std::sync::Arc<dyn Fn() -> Vec<PathBuf> + Send + Sync>,
startup_shell_overlay: zeph_core::ShellOverlaySnapshot,
shell_policy_handle: zeph_tools::ShellPolicyHandle,
#[cfg(feature = "scheduler")]
scheduler_executor: Option<std::sync::Arc<crate::scheduler_executor::SchedulerExecutor>>,
#[cfg(feature = "scheduler")]
scheduler_update_tx: Option<tokio::sync::broadcast::Sender<String>>,
#[cfg(feature = "scheduler")]
scheduler_custom_tx: Option<tokio::sync::broadcast::Sender<String>>,
}
#[cfg(feature = "acp")]
fn broadcast_to_mpsc<T: Clone + Send + 'static>(
mut brx: tokio::sync::broadcast::Receiver<T>,
cancel: zeph_memory::CancellationToken,
) -> tokio::sync::mpsc::Receiver<T> {
let (tx, rx) = tokio::sync::mpsc::channel(16);
tokio::spawn(async move {
loop {
tokio::select! {
() = cancel.cancelled() => break,
result = brx.recv() => {
match result {
Ok(item) => {
if tx.send(item).await.is_err() {
break; }
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(n)) => {
tracing::warn!(skipped = n, "broadcast_to_mpsc: lagged, some reload events dropped");
}
Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
}
}
}
}
});
rx
}
#[cfg(feature = "acp")]
pub(crate) struct PrebuiltAcpCore {
pub(crate) core: SharedCore,
pub(crate) supervisor: std::sync::Arc<zeph_common::TaskSupervisor>,
}
#[cfg(feature = "acp")]
#[allow(clippy::too_many_lines)]
async fn build_acp_deps(
app: &AppBuilder,
prebuilt_core: Option<PrebuiltAcpCore>,
prebuilt_mcp_manager: Option<std::sync::Arc<zeph_mcp::McpManager>>,
) -> anyhow::Result<(SharedAgentDeps, Box<dyn std::any::Any>)> {
log_acp_runtime_paths(app.config(), app.config_path());
let embed_model = app.embedding_model();
let (
SharedCore {
provider,
embedding_provider,
registry,
matcher,
memory,
budget_tokens,
rl_head,
},
acp_mem_supervisor,
) = if let Some(p) = prebuilt_core {
(p.core, p.supervisor)
} else {
let acp_mem_cancel = tokio_util::sync::CancellationToken::new();
let acp_mem_supervisor =
std::sync::Arc::new(zeph_common::TaskSupervisor::new(acp_mem_cancel));
let core = build_shared_core(app, &acp_mem_supervisor).await?;
(core, acp_mem_supervisor)
};
{
let sqlite = memory.sqlite().clone();
let retention_secs = app
.config()
.tools
.overflow
.retention_days
.saturating_mul(86_400);
let cell = std::sync::Arc::new(parking_lot::Mutex::new(Some((sqlite, retention_secs))));
acp_mem_supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "overflow_cleanup",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
let args = cell.lock().take();
async move {
if let Some((sqlite, retention_secs)) = args {
match sqlite.cleanup_overflow(retention_secs).await {
Ok(n) if n > 0 => {
tracing::info!("cleaned up {n} stale overflow entries");
}
Ok(_) => {}
Err(e) => tracing::warn!("overflow cleanup failed: {e}"),
}
} else {
tracing::warn!("overflow_cleanup factory called more than once");
}
}
},
});
}
let config = app.config();
if config.cli.safe_mode {
tracing::info!(
"safe mode active: ZEPH.md/CLAUDE.md/AGENTS.md, plugins, skills, hooks, and MCP \
servers are disabled for this session"
);
}
agent_setup::spawn_memory_maintenance_loops(
app,
&memory,
&provider,
&acp_mem_supervisor,
None,
false,
"acp",
);
let filter_registry = if config.tools.filters.enabled {
zeph_tools::OutputFilterRegistry::default_filters(&config.tools.filters)
} else {
zeph_tools::OutputFilterRegistry::new(false)
};
let permission_policy =
zeph_tools::build_permission_policy(&config.tools, config.security.autonomy_level);
let mut shell_executor = zeph_tools::ShellExecutor::new(&config.tools.shell)
.with_permissions(permission_policy.clone())
.with_output_filters(filter_registry)
.with_task_supervisor((*acp_mem_supervisor).clone());
if config.tools.sandbox.enabled {
let denied_present = !config.tools.sandbox.denied_domains.is_empty();
match zeph_tools::sandbox::build_sandbox_with_policy(
config.tools.sandbox.strict,
config.tools.sandbox.fail_if_unavailable,
denied_present,
) {
Ok(backend) => {
let name = backend.name();
let policy = crate::agent_setup::sandbox_policy_from_config(&config.tools.sandbox);
shell_executor = shell_executor.with_sandbox(std::sync::Arc::from(backend), policy);
tracing::info!(backend = name, "OS sandbox enabled (acp)");
}
Err(e) if config.tools.sandbox.strict || config.tools.sandbox.fail_if_unavailable => {
panic!("sandbox initialization failed: {e}");
}
Err(e) => {
tracing::warn!("OS sandbox unavailable, running without isolation: {e}");
}
}
}
let mut scrape_executor = zeph_tools::WebScrapeExecutor::new(&config.tools.scrape)
.with_egress_config(config.tools.egress.clone());
let web_search_api_key = config
.secrets
.web_search_api_key
.as_ref()
.map(|s| zeph_common::secret::Secret::new(s.expose()));
let mut web_search_executor = zeph_tools::WebSearchExecutor::new(
&config.tools.search,
&config.tools.scrape,
web_search_api_key,
)
.map(|w| w.with_egress_config(config.tools.egress.clone()));
if config.tools.egress.enabled {
let (egress_tx, egress_rx) = tokio::sync::mpsc::channel(256);
let dropped = std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0));
scrape_executor =
scrape_executor.with_egress_tx(egress_tx.clone(), std::sync::Arc::clone(&dropped));
if let Some(w) = web_search_executor.take() {
web_search_executor = Some(w.with_egress_tx(egress_tx, dropped));
}
{
let cell = std::sync::Arc::new(parking_lot::Mutex::new(Some(egress_rx)));
acp_mem_supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "egress_drain",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
let rx = cell.lock().take();
async move {
if let Some(rx) = rx {
agent_setup::drain_egress_events(rx, None).await;
} else {
tracing::warn!("egress_drain factory called more than once");
}
}
},
});
}
}
let mut acp_audit_logger: Option<std::sync::Arc<zeph_tools::AuditLogger>> = None;
if config.tools.audit.enabled
&& let Ok(logger) = zeph_tools::AuditLogger::from_config(&config.tools.audit, false).await
{
let logger = std::sync::Arc::new(logger);
shell_executor = shell_executor.with_audit(std::sync::Arc::clone(&logger));
scrape_executor = scrape_executor.with_audit(std::sync::Arc::clone(&logger));
if let Some(w) = web_search_executor.take() {
web_search_executor = Some(w.with_audit(std::sync::Arc::clone(&logger)));
}
acp_audit_logger = Some(logger);
}
let file_executor = zeph_tools::FileExecutor::new(
config
.tools
.shell
.allowed_paths
.iter()
.map(PathBuf::from)
.collect(),
);
let mcp_manager = if let Some(m) = prebuilt_mcp_manager {
m
} else {
let builder =
crate::bootstrap::create_mcp_manager_with_vault(config, false, app.age_vault_arc());
let builder =
crate::bootstrap::wire_trust_calibration(builder, config, Some(memory.sqlite().pool()))
.await;
std::sync::Arc::new(builder)
};
let (mcp_tools, _mcp_outcomes) = if config.cli.safe_mode {
(Vec::new(), Vec::new())
} else {
mcp_manager.connect_all().await
};
let mcp_shared_tools = std::sync::Arc::new(RwLock::new(mcp_tools.clone()));
let mut mcp_executor =
zeph_mcp::McpToolExecutor::new(mcp_manager.clone(), mcp_shared_tools.clone());
if config.cli.no_mcp_media {
tracing::info!("--no-mcp-media: MCP image passthrough disabled for this session");
} else {
mcp_executor = mcp_executor.with_media(
std::sync::Arc::new(zeph_sanitizer::MediaSanitizer::new(&config.mcp.media)),
config.mcp.media.max_images_per_result,
);
}
if let Some(ref logger) = acp_audit_logger {
mcp_executor = mcp_executor.with_audit(std::sync::Arc::clone(logger));
}
let shell_policy_handle = shell_executor.policy_handle();
let diagnostics_executor = crate::agent_setup::build_diagnostics_executor(config);
let clock: std::sync::Arc<dyn zeph_common::ClockSource> =
std::sync::Arc::new(zeph_common::SystemClock);
let time_executor = crate::agent_setup::build_time_executor(std::sync::Arc::clone(&clock));
let base_executor = crate::agent_setup::build_base_executor_chain(
file_executor,
shell_executor,
scrape_executor,
diagnostics_executor,
time_executor,
config
.tools
.shell
.allowed_paths
.iter()
.map(PathBuf::from)
.collect(),
);
let base_executor =
crate::agent_setup::with_search_executor(base_executor, web_search_executor);
let index_provider = crate::bootstrap::resolve_index_embed_provider(config, provider.clone());
let inner_executor: std::sync::Arc<dyn zeph_tools::ErasedToolExecutor> = {
let base: std::sync::Arc<dyn zeph_tools::ErasedToolExecutor> = std::sync::Arc::new(
zeph_tools::CompositeExecutor::new(base_executor, mcp_executor),
);
if let Some(search_executor) = crate::agent_setup::build_search_code_executor(
config,
app.qdrant_ops().cloned(),
index_provider.clone(),
memory.sqlite().pool().clone(),
Some(std::sync::Arc::clone(&mcp_manager)),
) {
std::sync::Arc::new(zeph_tools::CompositeExecutor::new(
zeph_tools::DynExecutor(base),
search_executor,
))
} else {
base
}
};
let tool_executor = inner_executor;
let policy_gate_pieces = agent_setup::build_policy_gate_pieces(config, &provider).await;
let capability_scopes_config = config.security.capability_scopes.clone();
let shadow_sentinel_config = config.security.shadow_sentinel.clone();
let shadow_sentinel_probe_provider = {
let sentinel_cfg = &shadow_sentinel_config;
let base = if sentinel_cfg.probe_provider.is_empty() {
provider.clone()
} else {
match crate::bootstrap::create_named_provider(
sentinel_cfg.probe_provider.as_str(),
config,
) {
Ok(p) => p,
Err(e) => {
tracing::warn!(
provider = %sentinel_cfg.probe_provider,
error = %e,
"shadow_sentinel probe provider resolution failed, using primary"
);
provider.clone()
}
}
};
match app.secret_registry() {
Some(registry) => {
base.masked(registry as std::sync::Arc<dyn zeph_llm::masking::OutboundMasker>)
}
None => base,
}
};
let trajectory_sentinel_config = config.security.trajectory.clone();
let quality_pipeline = crate::agent_setup::build_quality_pipeline(
config,
&provider,
app.secret_registry().as_ref(),
);
let mcp_registry = create_mcp_registry(
config,
&provider,
&mcp_tools,
&embed_model,
app.qdrant_ops(),
)
.await;
let summary_provider = app.build_summary_provider();
let skill_paths = app.skill_paths_for_registry();
let plugin_dirs_supplier = app.plugin_dirs_supplier();
let acp_project_rules = collect_project_rules(&skill_paths);
let crate::bootstrap::WatcherBundle {
skill_watcher,
skill_reload_rx: mpsc_skill_rx,
config_watcher,
config_reload_rx: mpsc_config_rx,
} = app.build_watchers(&acp_mem_supervisor);
let config_path_owned = app.config_path().to_owned();
let (_, shutdown_rx) = AppBuilder::build_shutdown();
let broadcast_cap = config.acp.broadcast_capacity.max(1);
let (skill_reload_tx, _) = tokio::sync::broadcast::channel(broadcast_cap);
let (config_reload_tx, _) = tokio::sync::broadcast::channel(broadcast_cap);
{
let skill_tx = skill_reload_tx.clone();
let cell = std::sync::Arc::new(parking_lot::Mutex::new(Some(mpsc_skill_rx)));
acp_mem_supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "skill_reload_fwd",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
let rx = cell.lock().take();
let tx = skill_tx.clone();
async move {
if let Some(mut rx) = rx {
while let Some(ev) = rx.recv().await {
let _ = tx.send(ev);
}
} else {
tracing::warn!("skill_reload_fwd factory called more than once");
}
}
},
});
}
{
let cfg_tx = config_reload_tx.clone();
let cell = std::sync::Arc::new(parking_lot::Mutex::new(Some(mpsc_config_rx)));
acp_mem_supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "config_reload_fwd",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
let rx = cell.lock().take();
let tx = cfg_tx.clone();
async move {
if let Some(mut rx) = rx {
while let Some(ev) = rx.recv().await {
let _ = tx.send(ev);
}
} else {
tracing::warn!("config_reload_fwd factory called more than once");
}
}
},
});
}
#[cfg(feature = "scheduler")]
let (scheduler_executor, scheduler_update_tx, scheduler_custom_tx) = {
let exp_deps = {
use std::sync::Arc;
if config.experiments.enabled && config.experiments.schedule.enabled {
let p = provider.clone();
let eval_provider = app.build_eval_provider().unwrap_or_else(|| p.clone());
Some((
Arc::new(p),
Arc::new(eval_provider),
Some(Arc::clone(&memory)),
))
} else {
None
}
};
let five_signal = memory.five_signal_runtime();
match crate::scheduler::init_scheduler(
config,
shutdown_rx.clone(),
exp_deps,
five_signal,
Some(&acp_mem_supervisor),
)
.await
{
Some(result) => {
let exec = std::sync::Arc::new(result.executor);
let custom_rx = result.custom_rx;
let (ctx, _) = tokio::sync::broadcast::channel::<String>(broadcast_cap);
let ctx_clone = ctx.clone();
let cell = std::sync::Arc::new(parking_lot::Mutex::new(Some(custom_rx)));
acp_mem_supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "sched_custom_fwd",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
let rx = cell.lock().take();
let tx = ctx_clone.clone();
async move {
if let Some(mut rx) = rx {
while let Some(ev) = rx.recv().await {
let _ = tx.send(ev);
}
} else {
tracing::warn!("sched_custom_fwd factory called more than once");
}
}
},
});
let update_tx = if let Some(update_rx) = result.update_rx {
let (utx, _) = tokio::sync::broadcast::channel::<String>(broadcast_cap);
let utx_clone = utx.clone();
let cell = std::sync::Arc::new(parking_lot::Mutex::new(Some(update_rx)));
acp_mem_supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "sched_update_fwd",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
let rx = cell.lock().take();
let tx = utx_clone.clone();
async move {
if let Some(mut rx) = rx {
while let Some(ev) = rx.recv().await {
let _ = tx.send(ev);
}
} else {
tracing::warn!(
"sched_update_fwd factory called more than once"
);
}
}
},
});
Some(utx)
} else {
None
};
let (update_tx, custom_tx) = (update_tx, Some(ctx));
(Some(exec), update_tx, custom_tx)
}
None => (None, None, None),
}
};
let session_config = zeph_core::AgentSessionConfig::from_config(config, budget_tokens);
let (resume_condenser_built, resume_token_counter_built) =
zeph_core::provider_factory::build_resume_condenser(config, &provider);
let feedback_classifier = app.build_feedback_classifier(&provider);
let provider_config_snapshot = agent_setup::build_provider_config_snapshot(config);
let acp_auth_clients = resolve_acp_auth_clients(&config.acp, app.vault()).await?;
let deps = SharedAgentDeps {
provider,
embedding_provider,
registry,
matcher,
max_active_skills: config.skills.max_active_skills.get(),
skill_disambiguation_threshold: config.skills.disambiguation_threshold,
skill_two_stage_matching: config.skills.two_stage_matching,
skill_confusability_threshold: config.skills.confusability_threshold,
skill_group_structured: config.skills.group_structured,
skill_support_similarity_threshold: config.skills.support_similarity_threshold,
skill_min_injection_score: config.skills.min_injection_score,
skill_generation_provider: config.skills.generation_provider.as_str().to_owned(),
skill_disambiguate_provider: config.skills.disambiguate_provider.as_str().to_owned(),
semantic_scan: config.skills.semantic_scan,
semantic_scan_provider: config.skills.semantic_scan_provider.as_str().to_owned(),
trust_config: config.skills.trust.clone(),
rl_routing_enabled: config.skills.rl_routing_enabled,
rl_learning_rate: config.skills.rl_learning_rate,
rl_weight: config.skills.rl_weight,
rl_persist_interval: config.skills.rl_persist_interval,
rl_warmup_updates: config.skills.rl_warmup_updates,
rl_head,
tool_executor,
clock,
permission_policy,
policy_gate_pieces,
capability_scopes_config,
shadow_sentinel_config,
shadow_sentinel_probe_provider,
trajectory_sentinel_config,
quality_pipeline,
skill_paths,
skill_reload_tx,
config_reload_tx,
memory,
history_limit: config.memory.history_limit,
recall_limit: config.memory.semantic.recall_limit,
summarization_threshold: config.memory.summarization_threshold,
shutdown_summary: config.memory.shutdown_summary,
shutdown_summary_min_messages: config.memory.shutdown_summary_min_messages,
shutdown_summary_max_messages: config.memory.shutdown_summary_max_messages,
shutdown_summary_timeout_secs: config.memory.shutdown_summary_timeout_secs,
shutdown_summary_provider: config.memory.shutdown_summary_provider.as_str().to_owned(),
channel_provider_persistence: config.session.provider_persistence,
channel_persist_provider_overrides: config.session.persist_provider_overrides,
index_config: config.index.clone(),
code_index_provider: index_provider,
code_qdrant_ops: app.qdrant_ops().cloned(),
shutdown_rx,
config_path: config_path_owned,
mcp_tools,
mcp_registry,
mcp_manager,
mcp_shared_tools,
mcp_config: config.mcp.clone(),
summary_provider,
judge_provider: app.build_judge_provider(),
feedback_classifier,
#[cfg(feature = "classifiers")]
classifiers_config: config.classifiers.clone(),
#[cfg(feature = "classifiers")]
pii_filter_enabled: config.security.pii_filter.enabled,
causal_ipi_config: config.security.causal_ipi.clone(),
causal_provider: config
.security
.causal_ipi
.provider
.as_deref()
.filter(|s| !s.is_empty())
.and_then(|name| match crate::bootstrap::create_named_provider(name, config) {
Ok(p) => {
tracing::info!(provider = %name, "causal IPI dedicated provider configured (acp)");
Some(p)
}
Err(e) => {
tracing::warn!(
provider = %name,
error = %e,
"causal IPI provider resolution failed, falling back to primary (acp)"
);
None
}
}),
nli_config: config.security.content_isolation.nli.clone(),
nli_provider: config
.security
.content_isolation
.nli
.provider
.as_non_empty()
.and_then(|name| match crate::bootstrap::create_named_provider(name, config) {
Ok(p) => {
tracing::info!(provider = %name, "NLI dedicated provider configured (acp)");
Some(p)
}
Err(e) => {
tracing::warn!(
provider = %name,
error = %e,
"NLI provider resolution failed, falling back to primary (acp)"
);
None
}
}),
secret_registry: app.secret_registry(),
vigil_config: config.security.vigil.clone(),
probe_provider: app.build_probe_provider(),
planner_provider: app.build_planner_provider(),
verify_provider: app.build_verify_provider(),
ensemble_members: app.build_ensemble_members(),
orchestrator_provider: app.build_orchestrator_provider(),
predicate_provider: app.build_predicate_provider(),
quarantine_provider: app.build_quarantine_provider(),
guardrail_provider: app.build_guardrail_provider(),
audit_logger: acp_audit_logger,
hooks_config: config.hooks.clone(),
safe_mode: config.cli.safe_mode,
cwd_allowed_paths: config
.tools
.shell
.allowed_paths
.iter()
.map(std::path::PathBuf::from)
.collect(),
tools_enabled: config.tools.enabled,
session_config,
session_persistence_config: config.session.clone(),
resume_condenser: resume_condenser_built,
resume_token_counter: resume_token_counter_built,
provider_pool: config.llm.providers.clone(),
provider_config_snapshot,
focus_config: config.agent.focus.clone(),
sidequest_config: config.memory.sidequest.clone(),
trajectory_config: config.memory.trajectory.clone(),
category_config: config.memory.category.clone(),
tool_filter_config: config.agent.tool_filter.clone(),
acp_agent_name: config.acp.agent_name.clone(),
acp_agent_version: config.acp.agent_version.clone(),
acp_max_sessions: config.acp.max_sessions,
acp_session_idle_timeout_secs: config.acp.session_idle_timeout_secs,
acp_permission_file: config.acp.permission_file.clone(),
acp_available_models: std::sync::Arc::new(RwLock::new(
if config.acp.available_models.is_empty() {
discover_models_from_config(config).await
} else {
config.acp.available_models.clone()
},
)),
acp_auth_clients,
acp_discovery_enabled: config.acp.discovery_enabled,
acp_title_max_chars: config.memory.sessions.title_max_chars,
acp_max_history: config.memory.sessions.max_history,
acp_log_file: if config.logging.file.is_empty() {
None
} else {
let cwd = std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."));
Some(
resolve_runtime_path(std::path::Path::new(&config.logging.file), &cwd)
.display()
.to_string(),
)
},
sqlite_path: crate::db_url::resolve_db_url(config).to_owned(),
acp_provider_factory: Some(build_acp_provider_factory(config, app.secret_registry())),
acp_provider_names: acp_provider_names(config),
acp_project_rules,
acp_additional_directories: config.acp.additional_directories.clone(),
acp_auth_methods: config.acp.auth_methods.clone(),
acp_message_ids_enabled: config.acp.message_ids_enabled,
acp_timeouts: config.acp.timeouts.clone(),
acp_model_config: config.acp.model_config.clone(),
plugin_dirs_supplier: std::sync::Arc::new(plugin_dirs_supplier),
#[cfg(feature = "scheduler")]
scheduler_executor,
#[cfg(feature = "scheduler")]
scheduler_update_tx,
#[cfg(feature = "scheduler")]
scheduler_custom_tx,
startup_shell_overlay: {
let mut blocked = config.tools.shell.blocked_commands.clone();
blocked.sort();
let mut allowed = config.tools.shell.allowed_commands.clone();
allowed.sort();
zeph_core::ShellOverlaySnapshot { blocked, allowed }
},
shell_policy_handle,
};
let keepalive: Box<dyn std::any::Any> = Box::new((skill_watcher, config_watcher));
Ok((deps, keepalive))
}
#[cfg(feature = "acp")]
const SESSION_LOCK_DEGRADED_MESSAGE: &str =
"Session persistence unavailable: another process already holds this session's write lock.";
#[cfg(feature = "acp")]
async fn notify_lock_degraded(
status_notifier: Option<&zeph_acp::SessionStatusNotifier>,
channel: &mut zeph_core::channel::LoopbackChannel,
) {
if let Some(notifier) = status_notifier {
notifier.notify_status_nowait(SESSION_LOCK_DEGRADED_MESSAGE);
} else {
channel
.send_status_best_effort(SESSION_LOCK_DEGRADED_MESSAGE)
.await;
}
}
#[cfg(feature = "acp")]
async fn open_session_log_or_notify_locked(
session_path: &std::path::Path,
status_notifier: Option<&zeph_acp::SessionStatusNotifier>,
channel: &mut zeph_core::channel::LoopbackChannel,
) -> Option<std::sync::Arc<zeph_session::SessionEventLog>> {
match zeph_session::SessionEventLog::open_exclusive(session_path).await {
Ok(log) => Some(std::sync::Arc::new(log)),
Err(zeph_session::SessionError::AlreadyLocked { path, pid, .. }) => {
tracing::error!(
lock_path = %path,
pid,
"failed to open session event log for ACP session: another process \
already holds this session's write lock; session persistence disabled \
for this session"
);
notify_lock_degraded(status_notifier, channel).await;
None
}
Err(e) => {
tracing::warn!(error = %e, "failed to open session event log for ACP session; session persistence disabled for this session");
None
}
}
}
#[cfg(feature = "acp")]
#[allow(clippy::struct_excessive_bools)]
struct BuildAcpAgentParams<F>
where
F: Fn() -> Vec<PathBuf> + Send + Sync + 'static,
{
provider: zeph_llm::any::AnyProvider,
embedding_provider: zeph_llm::any::AnyProvider,
registry: std::sync::Arc<RwLock<zeph_skills::registry::SkillRegistry>>,
matcher: Option<zeph_skills::matcher::SkillMatcherBackend>,
max_active_skills: usize,
tool_executor: zeph_tools::DynExecutor,
clock: std::sync::Arc<dyn zeph_common::ClockSource>,
session_config: zeph_core::AgentSessionConfig,
skill_disambiguation_threshold: f32,
skill_two_stage_matching: bool,
skill_confusability_threshold: f32,
skill_group_structured: bool,
skill_support_similarity_threshold: f32,
skill_min_injection_score: f32,
skill_generation_provider: String,
skill_disambiguate_provider: String,
semantic_scan: bool,
semantic_scan_provider: String,
trust_config: zeph_core::config::TrustConfig,
trust_snapshot:
std::sync::Arc<RwLock<std::collections::HashMap<String, zeph_core::SkillTrustSnapshot>>>,
quality_pipeline: Option<std::sync::Arc<zeph_core::quality::SelfCheckPipeline>>,
rl_routing_enabled: bool,
rl_learning_rate: f32,
rl_weight: f32,
rl_persist_interval: u32,
rl_warmup_updates: u32,
working_dir: PathBuf,
skill_paths: Vec<PathBuf>,
reload_rx: tokio::sync::mpsc::Receiver<zeph_skills::watcher::SkillEvent>,
plugin_dirs_supplier: F,
shutdown_rx: tokio::sync::watch::Receiver<bool>,
config_path: PathBuf,
config_reload_rx: tokio::sync::mpsc::Receiver<zeph_core::config_watcher::ConfigEvent>,
startup_shell_overlay: zeph_core::ShellOverlaySnapshot,
shell_policy_handle: zeph_tools::ShellPolicyHandle,
mcp_tools: Vec<zeph_mcp::McpTool>,
mcp_registry: Option<zeph_mcp::McpToolRegistry>,
mcp_manager: std::sync::Arc<zeph_mcp::McpManager>,
mcp_shared_tools: std::sync::Arc<RwLock<Vec<zeph_mcp::McpTool>>>,
mcp_config: zeph_core::config::McpConfig,
focus_config: zeph_core::config::FocusConfig,
sidequest_config: zeph_core::config::SidequestConfig,
trajectory_config: zeph_core::config::TrajectoryConfig,
category_config: zeph_core::config::CategoryConfig,
provider_pool: Vec<zeph_core::config::ProviderEntry>,
provider_config_snapshot: zeph_core::ProviderConfigSnapshot,
shutdown_summary: bool,
shutdown_summary_min_messages: usize,
shutdown_summary_max_messages: usize,
shutdown_summary_timeout_secs: u64,
shutdown_summary_provider: String,
channel_provider_persistence: bool,
channel_persist_provider_overrides: bool,
safe_mode: bool,
cwd_allowed_paths: Vec<PathBuf>,
tools_enabled: bool,
tool_filter_config: zeph_core::config::ToolFilterConfig,
}
#[cfg(feature = "acp")]
async fn build_acp_agent<C, F>(deps: BuildAcpAgentParams<F>, channel: C) -> Agent<C>
where
C: zeph_core::channel::Channel,
F: Fn() -> Vec<PathBuf> + Send + Sync + 'static,
{
Agent::new_with_registry_arc(
deps.provider.clone(),
deps.embedding_provider.clone(),
channel,
deps.registry,
deps.matcher,
deps.max_active_skills,
deps.tool_executor,
)
.apply_session_config(deps.session_config)
.with_skill_config(zeph_core::SkillConfigParams {
disambiguation_threshold: deps.skill_disambiguation_threshold,
two_stage_matching: deps.skill_two_stage_matching,
confusability_threshold: deps.skill_confusability_threshold,
group_structured: deps.skill_group_structured,
support_similarity_threshold: deps.skill_support_similarity_threshold,
min_injection_score: deps.skill_min_injection_score,
generation_provider_name: deps.skill_generation_provider,
disambiguate_provider_name: deps.skill_disambiguate_provider,
semantic_scan: deps.semantic_scan,
semantic_scan_provider_name: deps.semantic_scan_provider,
})
.with_trust_config(deps.trust_config)
.with_trust_snapshot(deps.trust_snapshot)
.with_quality_pipeline(deps.quality_pipeline)
.with_rl_routing(
deps.rl_routing_enabled,
deps.rl_learning_rate,
deps.rl_weight,
deps.rl_persist_interval,
deps.rl_warmup_updates,
)
.with_working_dir(deps.working_dir)
.with_skill_coldstart(
deps.skill_paths,
deps.reload_rx,
deps.plugin_dirs_supplier,
crate::bootstrap::managed_skills_dir(),
)
.with_shutdown(deps.shutdown_rx)
.with_config_reload(deps.config_path, deps.config_reload_rx)
.with_plugins_dir(crate::bootstrap::plugins_dir(), deps.startup_shell_overlay)
.with_shell_policy_handle(deps.shell_policy_handle)
.with_mcp(
deps.mcp_tools,
deps.mcp_registry,
Some(deps.mcp_manager),
&deps.mcp_config,
)
.with_mcp_shared_tools(deps.mcp_shared_tools)
.with_focus_and_sidequest_config(deps.focus_config, deps.sidequest_config)
.with_trajectory_and_category_config(deps.trajectory_config, deps.category_config)
.with_provider_pool(deps.provider_pool, deps.provider_config_snapshot)
.with_embedding_provider(deps.embedding_provider)
.with_shutdown_summary_config(
deps.shutdown_summary,
deps.shutdown_summary_min_messages,
deps.shutdown_summary_max_messages,
deps.shutdown_summary_timeout_secs,
)
.with_shutdown_summary_provider(deps.shutdown_summary_provider)
.with_channel_identity(
"acp",
deps.channel_provider_persistence,
deps.channel_persist_provider_overrides,
)
.with_safe_mode(deps.safe_mode)
.with_clock(deps.clock)
.with_allowed_paths(deps.cwd_allowed_paths)
.with_tools_enabled(deps.tools_enabled)
.maybe_init_tool_schema_filter(deps.tool_filter_config, deps.provider)
.await
}
#[cfg(feature = "acp")]
#[allow(clippy::too_many_lines)]
async fn spawn_acp_agent(
d: std::sync::Arc<SharedAgentDeps>,
mut channel: zeph_core::channel::LoopbackChannel,
acp_ctx: Option<zeph_acp::AcpContext>,
session_ctx: zeph_acp::SessionContext,
) {
use std::sync::Arc;
let provider = d.provider.clone();
let registry = Arc::clone(&d.registry);
let matcher = d.matcher.clone();
let max_active_skills = d.max_active_skills;
let skill_disambiguation_threshold = d.skill_disambiguation_threshold;
let skill_two_stage_matching = d.skill_two_stage_matching;
let skill_confusability_threshold = d.skill_confusability_threshold;
let skill_group_structured = d.skill_group_structured;
let skill_support_similarity_threshold = d.skill_support_similarity_threshold;
let skill_min_injection_score = d.skill_min_injection_score;
let skill_generation_provider = d.skill_generation_provider.clone();
let skill_disambiguate_provider = d.skill_disambiguate_provider.clone();
let semantic_scan = d.semantic_scan;
let semantic_scan_provider = d.semantic_scan_provider.clone();
let tool_executor = Arc::clone(&d.tool_executor);
let clock = Arc::clone(&d.clock);
let permission_policy = d.permission_policy.clone();
let skill_paths = d.skill_paths.clone();
let plugin_dirs_supplier = Arc::clone(&d.plugin_dirs_supplier);
let memory = Arc::clone(&d.memory);
let history_limit = d.history_limit;
let recall_limit = d.recall_limit;
let summarization_threshold = d.summarization_threshold;
let shutdown_summary = d.shutdown_summary;
let shutdown_summary_min_messages = d.shutdown_summary_min_messages;
let shutdown_summary_max_messages = d.shutdown_summary_max_messages;
let shutdown_summary_timeout_secs = d.shutdown_summary_timeout_secs;
let shutdown_summary_provider = d.shutdown_summary_provider.clone();
let channel_provider_persistence = d.channel_provider_persistence;
let channel_persist_provider_overrides = d.channel_persist_provider_overrides;
let index_config = d.index_config.clone();
let code_index_provider = d.code_index_provider.clone();
let code_qdrant_ops = d.code_qdrant_ops.clone();
let shutdown_rx = d.shutdown_rx.clone();
let config_path = d.config_path.clone();
let mcp_tools = d.mcp_tools.clone();
let mcp_registry = d.mcp_registry.clone();
let mcp_manager = Arc::clone(&d.mcp_manager);
let mcp_shared_tools = Arc::clone(&d.mcp_shared_tools);
let mcp_config = d.mcp_config.clone();
let summary_provider = d.summary_provider.clone();
let judge_provider = d.judge_provider.clone();
let feedback_classifier = d.feedback_classifier.clone();
#[cfg(feature = "classifiers")]
let classifiers_config = d.classifiers_config.clone();
#[cfg(feature = "classifiers")]
let pii_filter_enabled = d.pii_filter_enabled;
let causal_ipi_config = d.causal_ipi_config.clone();
let causal_provider = d.causal_provider.clone();
let nli_config = d.nli_config.clone();
let nli_provider = d.nli_provider.clone();
let secret_registry = d.secret_registry.clone();
let vigil_config = d.vigil_config.clone();
let probe_provider = d.probe_provider.clone();
let planner_provider = d.planner_provider.clone();
let verify_provider = d.verify_provider.clone();
let ensemble_members = d.ensemble_members.clone();
let orchestrator_provider = d.orchestrator_provider.clone();
let predicate_provider = d.predicate_provider.clone();
let quarantine_provider = d.quarantine_provider.clone();
let guardrail_provider = d.guardrail_provider.clone();
let session_config = d.session_config.clone();
let session_persistence_config = d.session_persistence_config.clone();
let provider_pool = d.provider_pool.clone();
let provider_config_snapshot = d.provider_config_snapshot.clone();
let skill_reload_tx = d.skill_reload_tx.clone();
let config_reload_tx = d.config_reload_tx.clone();
#[cfg(feature = "scheduler")]
let scheduler_executor = d.scheduler_executor.as_ref().map(std::sync::Arc::clone);
#[cfg(feature = "scheduler")]
let scheduler_update_tx = d.scheduler_update_tx.clone();
#[cfg(feature = "scheduler")]
let scheduler_custom_tx = d.scheduler_custom_tx.clone();
let hooks_config = d.hooks_config.clone();
let safe_mode = d.safe_mode;
let cwd_allowed_paths = d.cwd_allowed_paths.clone();
let tools_enabled = d.tools_enabled;
let tool_filter_config = d.tool_filter_config.clone();
let status_notifier = acp_ctx.as_ref().map(|ctx| ctx.status_notifier.clone());
let adapter_cancel = zeph_memory::CancellationToken::new();
let reload_rx = broadcast_to_mpsc(skill_reload_tx.subscribe(), adapter_cancel.clone());
let config_reload_rx = broadcast_to_mpsc(config_reload_tx.subscribe(), adapter_cancel.clone());
#[cfg(feature = "scheduler")]
let scheduler_update_rx = scheduler_update_tx
.as_ref()
.map(|tx| broadcast_to_mpsc(tx.subscribe(), adapter_cancel.clone()));
#[cfg(feature = "scheduler")]
let scheduler_custom_rx = scheduler_custom_tx
.as_ref()
.map(|tx| broadcast_to_mpsc(tx.subscribe(), adapter_cancel.clone()));
let debug_config = session_config.debug_config.clone();
let memory_validation_config = session_config.security.memory_validation.clone();
let memory_executor = zeph_core::memory_tools::MemoryToolExecutor::with_validator(
Arc::clone(&memory),
session_ctx
.conversation_id
.unwrap_or(zeph_memory::ConversationId(0)),
zeph_sanitizer::memory_validation::MemoryWriteValidator::new(memory_validation_config),
);
let overflow_executor = {
let mut ex =
zeph_core::overflow_tools::OverflowToolExecutor::new(Arc::new(memory.sqlite().clone()));
if let Some(cid) = session_ctx.conversation_id {
ex = ex.with_conversation(cid.0);
}
ex
};
let (skill_loader_executor, skill_invoke_executor, trust_snapshot) =
agent_setup::build_skill_executors(®istry);
let (base_composite, cancel_signal, provider_override, parent_tool_use_id): (
Arc<dyn ErasedToolExecutor>,
_,
_,
_,
) = if let Some(ctx) = acp_ctx {
let cancel_signal = Arc::clone(&ctx.cancel_signal);
let provider_override = Arc::clone(&ctx.provider_override);
let parent_tool_use_id = ctx.parent_tool_use_id.clone();
let adapter_cancel_clone = adapter_cancel.clone();
let cancel_signal_clone = Arc::clone(&cancel_signal);
tokio::spawn(async move {
cancel_signal_clone.notified().await;
adapter_cancel_clone.cancel();
});
let mut base: Arc<dyn ErasedToolExecutor> = Arc::clone(&tool_executor) as Arc<_>;
if let Some(fs) = ctx.file_executor {
let filtered = zeph_tools::ToolFilter::new(
zeph_tools::DynExecutor(base),
&["read", "write", "glob"],
);
base = Arc::new(zeph_tools::CompositeExecutor::new(fs, filtered));
}
if let Some(shell) = ctx.shell_executor {
base = Arc::new(zeph_tools::CompositeExecutor::new(
shell,
zeph_tools::DynExecutor(base),
));
}
base = Arc::new(zeph_tools::CompositeExecutor::new(
skill_loader_executor,
zeph_tools::CompositeExecutor::new(
skill_invoke_executor,
zeph_tools::CompositeExecutor::new(
memory_executor,
zeph_tools::CompositeExecutor::new(
overflow_executor,
zeph_tools::DynExecutor(base),
),
),
),
));
(
base,
Some(cancel_signal),
Some(provider_override),
parent_tool_use_id,
)
} else {
let base: Arc<dyn ErasedToolExecutor> = Arc::new(zeph_tools::CompositeExecutor::new(
skill_loader_executor,
zeph_tools::CompositeExecutor::new(
skill_invoke_executor,
zeph_tools::CompositeExecutor::new(
memory_executor,
zeph_tools::CompositeExecutor::new(
overflow_executor,
zeph_tools::DynExecutor(Arc::clone(&tool_executor) as Arc<_>),
),
),
),
));
(base, None, None, None)
};
let (trust_gated, mcp_ids_handle) = crate::agent_setup::apply_common_tool_gating(
zeph_tools::DynExecutor(base_composite),
&permission_policy,
);
crate::agent_setup::register_mcp_tool_ids(&mcp_ids_handle, &mcp_tools);
let trajectory_risk_slot: zeph_tools::TrajectoryRiskSlot =
Arc::new(parking_lot::RwLock::new(0u8));
let trajectory_signal_queue: zeph_tools::RiskSignalQueue =
Arc::new(parking_lot::Mutex::new(Vec::new()));
let tool_executor = crate::agent_setup::apply_policy_gate_chain(
trust_gated,
&d.policy_gate_pieces,
d.audit_logger.as_ref(),
Some((&trajectory_risk_slot, &trajectory_signal_queue)),
);
let tool_executor = {
let scopes_cfg = &d.capability_scopes_config;
if scopes_cfg.scopes.is_empty() {
tool_executor
} else {
use std::collections::HashSet;
use zeph_tools::executor::ToolExecutor as _;
use zeph_tools::scope::build_scoped_executor;
let registry_ids: HashSet<String> = tool_executor
.tool_definitions()
.into_iter()
.map(|def| {
let id = def.id.to_string();
if id.contains(':') {
id
} else {
format!("builtin:{id}")
}
})
.collect();
let fallback = zeph_tools::DynExecutor(Arc::clone(&tool_executor.0));
match build_scoped_executor(tool_executor, scopes_cfg, ®istry_ids) {
Ok(scoped) => {
let scoped = scoped.with_signal_queue(Arc::clone(&trajectory_signal_queue));
zeph_tools::DynExecutor(Arc::new(scoped))
}
Err(e) => {
tracing::error!(
"capability_scopes: {e}, denying all tool access for this session \
(fail-closed)"
);
zeph_tools::DynExecutor(Arc::new(zeph_tools::scope::ScopedToolExecutor::new(
fallback,
zeph_tools::scope::ToolScope::empty(),
)))
}
}
}
};
let (tool_executor, shadow_sentinel_arc) = {
let sentinel_cfg = &d.shadow_sentinel_config;
if sentinel_cfg.enabled {
let pool = memory.sqlite().pool().clone();
let llm_probe = zeph_core::agent::shadow_sentinel::LlmSafetyProbe::new(
Arc::new(d.shadow_sentinel_probe_provider.clone()),
sentinel_cfg.probe_timeout_ms,
sentinel_cfg.deny_on_timeout,
);
let store = zeph_core::agent::shadow_sentinel::ShadowEventStore::new(pool);
let conversation_identity = session_ctx
.conversation_id
.unwrap_or(zeph_memory::ConversationId(0))
.0
.to_string();
let sentinel = Arc::new(zeph_core::agent::shadow_sentinel::ShadowSentinel::new(
store,
Box::new(llm_probe),
sentinel_cfg.clone(),
conversation_identity,
));
let turn_number = Arc::new(std::sync::atomic::AtomicU64::new(0));
let risk_level = Arc::new(parking_lot::RwLock::new("calm".to_owned()));
let probe_gate: Arc<dyn zeph_tools::ProbeGate> =
Arc::new(crate::runner::ShadowSentinelProbeGateAdapter {
sentinel: Arc::clone(&sentinel),
});
let shadow_exec = zeph_tools::ShadowProbeExecutor::new(
tool_executor,
probe_gate,
turn_number,
risk_level,
);
tracing::info!("security.shadow_sentinel: ShadowProbeExecutor wired (acp session)");
(
zeph_tools::DynExecutor(Arc::new(shadow_exec)),
Some(sentinel),
)
} else {
(tool_executor, None)
}
};
if let Some(ref sentinel) = shadow_sentinel_arc {
crate::agent_setup::register_mcp_tool_ids(&sentinel.mcp_tool_ids_handle(), &mcp_tools);
}
let mut acp_session_sink = None;
let mut preloaded_messages: Vec<zeph_llm::provider::Message> = Vec::new();
if session_persistence_config.enabled {
let sid = zeph_common::SessionId::new(session_ctx.session_id.to_string());
let store = zeph_session::SessionStore::new(memory.sqlite().pool().clone());
if let Err(e) = store.create(sid.as_str()).await {
tracing::warn!(error = %e, session_id = %sid, "failed to create session-store row for ACP session");
}
let data_dir = std::path::PathBuf::from(&session_persistence_config.data_dir);
let session_path = zeph_session::session_dir(&data_dir, sid.as_str());
let log = if let Some(cid) = session_ctx.conversation_id {
match zeph_agent_persistence::hydrate_and_condense(
&session_path,
&store,
sid.as_str(),
cid,
&memory,
None,
&d.resume_condenser,
d.resume_token_counter.as_ref(),
d.session_config.budget_tokens,
)
.await
{
Ok(hydrated) => {
preloaded_messages = hydrated.messages;
Some(hydrated.log)
}
Err(zeph_agent_persistence::PersistenceError::Session(
zeph_session::SessionError::AlreadyLocked { path, pid, .. },
)) => {
tracing::error!(
lock_path = %path,
pid,
"session hydration failed: another process already holds this session's \
write lock; session persistence disabled for this session"
);
notify_lock_degraded(status_notifier.as_ref(), &mut channel).await;
None
}
Err(e) => {
tracing::warn!(error = %e, "session hydration failed; session persistence disabled for this session");
None
}
}
} else {
open_session_log_or_notify_locked(&session_path, status_notifier.as_ref(), &mut channel)
.await
};
if let Some(log) = log {
acp_session_sink = Some(Arc::new(zeph_agent_persistence::SessionSink::new(
log, store, sid,
)));
}
}
let build_params = BuildAcpAgentParams {
provider: provider.clone(),
embedding_provider: d.embedding_provider.clone(),
registry: Arc::clone(®istry),
matcher,
max_active_skills,
tool_executor,
clock,
session_config,
skill_disambiguation_threshold,
skill_two_stage_matching,
skill_confusability_threshold,
skill_group_structured,
skill_support_similarity_threshold,
skill_min_injection_score,
skill_generation_provider,
skill_disambiguate_provider,
semantic_scan,
semantic_scan_provider,
trust_config: d.trust_config.clone(),
trust_snapshot: Arc::clone(&trust_snapshot),
quality_pipeline: d.quality_pipeline.clone(),
rl_routing_enabled: d.rl_routing_enabled,
rl_learning_rate: d.rl_learning_rate,
rl_weight: d.rl_weight,
rl_persist_interval: d.rl_persist_interval,
rl_warmup_updates: d.rl_warmup_updates,
working_dir: session_ctx.working_dir.clone(),
skill_paths,
reload_rx,
plugin_dirs_supplier: move || plugin_dirs_supplier(),
shutdown_rx,
config_path,
config_reload_rx,
startup_shell_overlay: d.startup_shell_overlay.clone(),
shell_policy_handle: d.shell_policy_handle.clone(),
mcp_tools,
mcp_registry,
mcp_manager: Arc::clone(&mcp_manager),
mcp_shared_tools,
mcp_config,
focus_config: d.focus_config.clone(),
sidequest_config: d.sidequest_config.clone(),
trajectory_config: d.trajectory_config.clone(),
category_config: d.category_config.clone(),
provider_pool,
provider_config_snapshot,
shutdown_summary,
shutdown_summary_min_messages,
shutdown_summary_max_messages,
shutdown_summary_timeout_secs,
shutdown_summary_provider,
channel_provider_persistence,
channel_persist_provider_overrides,
safe_mode,
cwd_allowed_paths,
tools_enabled,
tool_filter_config,
};
let mut agent = Box::pin(build_acp_agent(build_params, channel)).await;
agent = agent.with_acp_session(true);
agent = agent_setup::apply_code_retrieval(agent, &index_config);
agent = agent_setup::apply_code_rag_retriever(
agent,
&index_config,
code_qdrant_ops,
code_index_provider,
memory.sqlite().pool().clone(),
);
agent = agent
.with_trajectory_risk_slot(trajectory_risk_slot)
.with_signal_queue(trajectory_signal_queue)
.with_trajectory_config(d.trajectory_sentinel_config.clone())
.0;
if let Some(head) = d.rl_head.clone() {
agent = agent.with_rl_head(head);
}
if let Some(ref logger) = d.audit_logger {
agent = agent.with_audit_logger(std::sync::Arc::clone(logger));
}
#[cfg(feature = "scheduler")]
{
if let Some(rx) = scheduler_update_rx {
agent = agent.with_update_notifications(rx);
}
if let Some(rx) = scheduler_custom_rx {
agent = agent.with_custom_task_rx(rx);
}
if let Some(sched_exec) = scheduler_executor {
agent = agent.add_tool_executor(zeph_tools::DynExecutor(sched_exec));
}
}
if let Some(cid) = session_ctx.conversation_id {
agent = agent.with_memory(
Arc::clone(&memory),
cid,
history_limit,
recall_limit,
summarization_threshold,
);
}
if !preloaded_messages.is_empty() {
agent = agent.with_preloaded_messages(preloaded_messages);
}
if let Some(sink) = acp_session_sink {
agent = agent
.with_session_sink(Some(sink))
.with_session_persistence_config(Some(session_persistence_config.clone()));
}
if let Some(signal) = cancel_signal {
agent = agent.with_cancel_signal(signal);
}
if let Some(slot) = provider_override {
agent = agent.with_provider_override(slot);
}
if let Some(parent_id) = parent_tool_use_id {
agent = agent.with_parent_tool_use_id(parent_id);
}
if let Some(sp) = summary_provider {
agent = agent.with_summary_provider(sp);
}
if let Some(jp) = judge_provider {
agent = agent.with_judge_provider(jp);
}
if let Some(fc) = feedback_classifier {
agent = agent.with_llm_classifier(fc);
}
if let Some(pp) = probe_provider {
agent = agent.with_probe_provider(pp);
}
if let Some(pp) = planner_provider {
agent = agent.with_planner_provider(pp);
}
if let Some(vp) = verify_provider {
agent = agent.with_verify_provider(vp);
}
agent = agent.with_ensemble_members(ensemble_members);
if let Some(op) = orchestrator_provider {
agent = agent.with_orchestrator_provider(op);
}
if let Some(pp) = predicate_provider {
agent = agent.with_predicate_provider(pp);
}
agent = agent_setup::apply_quarantine_provider(agent, quarantine_provider);
{
agent = agent_setup::apply_guardrail(agent, guardrail_provider);
}
#[cfg(feature = "classifiers")]
{
agent = agent_setup::apply_injection_classifier_with_cfg(agent, &classifiers_config);
if classifiers_config.enabled {
agent = agent.with_enforcement_mode(classifiers_config.enforcement_mode);
}
agent = agent_setup::apply_three_class_classifier_with_cfg(agent, &classifiers_config);
agent = agent_setup::apply_pii_classifier_with_cfg(agent, &classifiers_config);
agent = agent_setup::apply_pii_ner_classifier_with_cfg(
agent,
&classifiers_config,
pii_filter_enabled,
);
}
agent = agent_setup::apply_causal_analyzer_with_cfg(
agent,
provider.clone(),
causal_provider,
&causal_ipi_config,
secret_registry.as_ref(),
);
agent = agent_setup::apply_nli_sanitizer_with_cfg(
agent,
provider.clone(),
nli_provider,
&nli_config,
secret_registry.as_ref(),
);
agent = agent_setup::apply_secret_masking(agent, secret_registry);
agent = agent_setup::apply_vigil(agent, &vigil_config);
if debug_config.enabled {
let session_dump_dir = debug_config
.output_dir
.join(session_ctx.session_id.to_string());
agent = agent_setup::apply_debug_dumper(
agent,
session_dump_dir.as_path(),
debug_config.format,
debug_config.include_raw_images,
)
.0;
}
if !safe_mode {
agent = agent.with_hooks_config(&hooks_config);
}
if let Some(sentinel) = shadow_sentinel_arc {
agent = agent.with_shadow_sentinel(sentinel);
}
agent = agent.with_mcp_tool_ids_handle(mcp_ids_handle);
drop(d);
if let Err(e) = agent.load_history().await {
tracing::error!("failed to load agent history: {e:#}");
}
if let Err(e) = Box::pin(agent.run()).await {
tracing::error!("ACP agent loop error: {e:#}");
}
agent.shutdown().await;
adapter_cancel.cancel();
}
#[cfg(feature = "acp")]
async fn discover_models_from_config(config: &zeph_core::config::Config) -> Vec<String> {
use zeph_llm::model_cache::ModelCache;
async fn expand_from_cache(slug: &str, fallback: &str) -> Vec<String> {
let cache = ModelCache::for_slug(slug);
if !cache.is_stale_async().await
&& let Ok(Some(entries)) = cache.load_async().await
&& !entries.is_empty()
{
return entries
.into_iter()
.map(|m| format!("{slug}:{}", m.id))
.collect();
}
vec![format!("{slug}:{fallback}")]
}
let mut models: Vec<String> = Vec::new();
for entry in &config.llm.providers {
let slug = entry.provider_type.as_str();
let fallback = entry.model.as_deref().unwrap_or("unknown");
models.extend(expand_from_cache(slug, fallback).await);
}
models.dedup();
models
}
#[cfg(feature = "acp")]
#[allow(clippy::too_many_lines)]
fn build_acp_provider_factory(
config: &zeph_core::config::Config,
secret_registry: Option<std::sync::Arc<zeph_sanitizer::secret_mask::SecretMaskRegistry>>,
) -> zeph_acp::ProviderFactory {
#[derive(Clone)]
enum ProviderSnapshot {
Ollama {
base_url: String,
embed: String,
},
Claude {
api_key: String,
max_tokens: u32,
},
OpenAi {
api_key: String,
base_url: String,
max_tokens: u32,
embed: Option<String>,
reasoning_effort: Option<String>,
},
Compatible {
api_key: String,
base_url: String,
max_tokens: u32,
embed: Option<String>,
name: String,
},
}
let mut snapshots: Vec<ProviderSnapshot> = Vec::new();
for entry in &config.llm.providers {
let name = entry.effective_name();
match entry.provider_type {
zeph_core::config::ProviderKind::Ollama => {
snapshots.push(ProviderSnapshot::Ollama {
base_url: entry
.base_url
.clone()
.unwrap_or_else(|| "http://localhost:11434".to_owned()),
embed: config.llm.embedding_model.clone(),
});
}
zeph_core::config::ProviderKind::Claude => {
if let Some(ref secret) = config.secrets.claude_api_key {
snapshots.push(ProviderSnapshot::Claude {
api_key: secret.expose().to_owned(),
max_tokens: entry.max_tokens.unwrap_or(4096),
});
}
}
zeph_core::config::ProviderKind::OpenAi => {
if let Some(ref secret) = config.secrets.openai_api_key {
snapshots.push(ProviderSnapshot::OpenAi {
api_key: secret.expose().to_owned(),
base_url: entry
.base_url
.clone()
.unwrap_or_else(|| "https://api.openai.com/v1".to_owned()),
max_tokens: entry.max_tokens.unwrap_or(4096),
embed: entry.embedding_model.clone(),
reasoning_effort: entry.reasoning_effort.clone(),
});
}
}
zeph_core::config::ProviderKind::Compatible => {
let secret = entry
.api_key
.as_deref()
.map(std::borrow::ToOwned::to_owned)
.or_else(|| {
config
.secrets
.compatible_api_keys
.get(&name)
.map(|s| s.expose().to_owned())
});
if let Some(api_key) = secret {
snapshots.push(ProviderSnapshot::Compatible {
api_key,
base_url: entry.base_url.clone().unwrap_or_default(),
max_tokens: entry.max_tokens.unwrap_or(4096),
embed: entry.embedding_model.clone(),
name,
});
}
}
_ => {}
}
}
let masker: Option<std::sync::Arc<dyn zeph_llm::masking::OutboundMasker>> =
secret_registry.map(|r| r as std::sync::Arc<dyn zeph_llm::masking::OutboundMasker>);
let snapshots = std::sync::Arc::new(snapshots);
std::sync::Arc::new(move |key: &str| {
let wrap = |p: zeph_llm::any::AnyProvider| -> zeph_llm::any::AnyProvider {
match &masker {
Some(m) => p.masked(std::sync::Arc::clone(m)),
None => p,
}
};
let (provider_name, model) = key.split_once(':')?;
let model = model.to_owned();
for snapshot in snapshots.as_ref() {
match snapshot {
ProviderSnapshot::Ollama {
base_url, embed, ..
} if provider_name == "ollama" => {
let mut p = zeph_llm::ollama::OllamaProvider::new(
base_url,
model.clone(),
embed.clone(),
);
p.set_context_window(0);
return Some(wrap(zeph_llm::any::AnyProvider::Ollama(p)));
}
ProviderSnapshot::Claude {
api_key,
max_tokens,
} if provider_name == "claude" => {
return Some(wrap(zeph_llm::any::AnyProvider::Claude(
zeph_llm::claude::ClaudeProvider::new(
api_key.clone(),
model.clone(),
*max_tokens,
),
)));
}
ProviderSnapshot::OpenAi {
api_key,
base_url,
max_tokens,
embed,
reasoning_effort,
} if provider_name == "openai" => {
return Some(wrap(zeph_llm::any::AnyProvider::OpenAi(
zeph_llm::openai::OpenAiProvider::new(zeph_llm::openai::OpenAiConfig {
api_key: api_key.clone(),
base_url: base_url.clone(),
model: model.clone(),
max_tokens: *max_tokens,
embedding_model: embed.clone(),
reasoning_effort: reasoning_effort.clone(),
context_window: None,
completion_tokens_param: None,
vision: None,
}),
)));
}
ProviderSnapshot::Compatible {
api_key,
base_url,
max_tokens,
embed,
name,
} if provider_name == name => {
return Some(wrap(zeph_llm::any::AnyProvider::Compatible(
zeph_llm::compatible::CompatibleProvider::new(
zeph_llm::compatible::CompatibleConfig {
provider_name: name.clone(),
api_key: api_key.clone(),
base_url: base_url.clone(),
model: model.clone(),
max_tokens: *max_tokens,
embedding_model: embed.clone(),
completion_tokens_param: None,
vision: None,
},
),
)));
}
_ => {}
}
}
None
})
}
#[cfg(feature = "acp")]
fn acp_provider_names(config: &zeph_core::config::Config) -> Vec<(String, zeph_acp::LlmProtocol)> {
config
.llm
.providers
.iter()
.map(|entry| {
let protocol = match entry.provider_type {
zeph_core::config::ProviderKind::Claude => zeph_acp::LlmProtocol::Anthropic,
zeph_core::config::ProviderKind::OpenAi
| zeph_core::config::ProviderKind::Compatible => zeph_acp::LlmProtocol::OpenAi,
other => zeph_acp::LlmProtocol::Other(other.as_str().to_owned()),
};
(entry.effective_name(), protocol)
})
.collect()
}
#[cfg(feature = "acp")]
fn collect_project_rules(skill_paths: &[PathBuf]) -> Vec<PathBuf> {
let mut rules = Vec::new();
let rules_dir = std::path::Path::new(".claude/rules");
if rules_dir.is_dir()
&& let Ok(entries) = std::fs::read_dir(rules_dir)
{
let mut paths: Vec<PathBuf> = entries
.flatten()
.map(|e| e.path())
.filter(|p| p.extension().is_some_and(|e| e == "md"))
.collect();
paths.sort();
rules.extend(paths);
}
for sp in skill_paths {
if sp.is_file() {
rules.push(sp.clone());
}
}
rules
}
#[cfg(feature = "acp")]
#[allow(clippy::too_many_arguments)] pub(crate) async fn run_acp_server(
config_path: Option<&std::path::Path>,
vault_backend: Option<&str>,
vault_key: Option<&std::path::Path>,
vault_path: Option<&std::path::Path>,
cli_additional_dirs: Vec<std::path::PathBuf>,
cli_auth_methods: Vec<String>,
cli_message_ids: Option<bool>,
safe_mode: bool,
no_mcp_media: bool,
) -> anyhow::Result<()> {
use std::sync::Arc;
let app = AppBuilder::new(
config_path,
vault_backend,
vault_key,
vault_path,
safe_mode,
no_mcp_media,
)
.await?;
let (mut deps, _keepalive) = Box::pin(build_acp_deps(&app, None, None)).await?;
let available_models = std::sync::Arc::clone(&deps.acp_available_models);
let provider = deps.provider.clone();
zeph_acp::warm_model_caches(provider, available_models).await;
let effective_additional_dirs = if cli_additional_dirs.is_empty() {
deps.acp_additional_directories.clone()
} else {
cli_additional_dirs
.into_iter()
.map(|p| {
zeph_core::config::AdditionalDir::parse(p.clone()).map_err(|e| {
anyhow::anyhow!("invalid --acp-additional-dir {}: {e}", p.display())
})
})
.collect::<anyhow::Result<Vec<_>>>()?
};
let effective_auth_methods = if cli_auth_methods.is_empty() {
let methods = deps.acp_auth_methods.clone();
anyhow::ensure!(
!methods.is_empty(),
"acp.auth_methods must not be empty; set at least one method (e.g. \"agent\")"
);
methods
} else {
let methods: Vec<_> = cli_auth_methods
.iter()
.map(|m| match m.as_str() {
"agent" => Ok(zeph_core::config::AcpAuthMethod::Agent),
other => Err(anyhow::anyhow!(
"unknown --acp-auth-method {other:?}; accepted values: agent"
)),
})
.collect::<anyhow::Result<Vec<_>>>()?;
anyhow::ensure!(
!methods.is_empty(),
"--acp-auth-method list must not be empty after parsing"
);
methods
};
let effective_message_ids = cli_message_ids.unwrap_or(deps.acp_message_ids_enabled);
let mcp_manager_for_acp = Arc::clone(&deps.mcp_manager);
let server_config = zeph_acp::AcpServerConfig {
agent_name: deps.acp_agent_name.clone(),
agent_version: deps.acp_agent_version.clone(),
max_sessions: deps.acp_max_sessions,
session_idle_timeout_secs: deps.acp_session_idle_timeout_secs,
permission_file: deps.acp_permission_file.clone(),
provider_factory: deps.acp_provider_factory.take(),
available_models: std::sync::Arc::clone(&deps.acp_available_models),
provider_names: deps.acp_provider_names.clone(),
mcp_manager: Some(mcp_manager_for_acp),
auth_clients: deps.acp_auth_clients.clone(),
discovery_enabled: deps.acp_discovery_enabled,
terminal_timeout_secs: deps.acp_timeouts.terminal_secs,
project_rules: deps.acp_project_rules.clone(),
title_max_chars: deps.acp_title_max_chars,
max_history: deps.acp_max_history,
sqlite_path: Some(deps.sqlite_path.clone()),
session_data_dir: deps
.session_persistence_config
.enabled
.then(|| std::path::PathBuf::from(&deps.session_persistence_config.data_dir)),
ready_notification: Some(zeph_acp::transport::ReadyNotification {
version: deps.acp_agent_version.clone(),
pid: std::process::id(),
log_file: deps.acp_log_file.clone(),
}),
additional_directories: effective_additional_dirs,
auth_methods: effective_auth_methods,
message_ids_enabled: effective_message_ids,
timeouts: deps.acp_timeouts.clone(),
model_config: deps.acp_model_config.clone(),
};
let shared = Arc::new(deps);
let spawner: zeph_acp::AgentSpawner = Arc::new(move |channel, acp_ctx, session_ctx| {
let shared = Arc::clone(&shared);
Box::pin(spawn_acp_agent(shared, channel, acp_ctx, session_ctx))
});
zeph_acp::serve_stdio(spawner, server_config).await?;
Ok(())
}
#[cfg(feature = "acp-http")]
#[allow(clippy::too_many_lines, clippy::too_many_arguments)] pub(crate) async fn run_acp_http_server(
config_path: Option<&std::path::Path>,
vault_backend: Option<&str>,
vault_key: Option<&std::path::Path>,
vault_path: Option<&std::path::Path>,
bind_override: Option<&str>,
auth_token_override: Option<String>,
safe_mode: bool,
no_mcp_media: bool,
) -> anyhow::Result<()> {
use std::sync::Arc;
use tokio::sync::RwLock;
let app = AppBuilder::new(
config_path,
vault_backend,
vault_key,
vault_path,
safe_mode,
no_mcp_media,
)
.await?;
log_acp_runtime_paths(app.config(), app.config_path());
let bind_addr = bind_override.map_or_else(|| app.config().acp.http_bind.clone(), str::to_owned);
let mut auth_clients = resolve_acp_auth_clients(&app.config().acp, app.vault()).await?;
if let Some(override_token) = auth_token_override {
auth_clients.retain(|c| c.id != zeph_config::ACP_AUTH_CLIENT_ID_DEFAULT);
anyhow::ensure!(
!auth_clients.iter().any(|c| c.token == override_token),
"--acp-auth-token collides with a configured [[acp.auth_clients]] token"
);
auth_clients.insert(
0,
zeph_acp::AcpClientToken {
id: zeph_config::ACP_AUTH_CLIENT_ID_DEFAULT.to_owned(),
token: override_token,
},
);
}
let mcp_manager_for_acp = Arc::new(crate::bootstrap::create_mcp_manager_with_vault(
app.config(),
false,
app.age_vault_arc(),
));
let server_config = zeph_acp::AcpServerConfig {
agent_name: app.config().acp.agent_name.clone(),
agent_version: app.config().acp.agent_version.clone(),
max_sessions: app.config().acp.max_sessions,
session_idle_timeout_secs: app.config().acp.session_idle_timeout_secs,
permission_file: app.config().acp.permission_file.clone(),
provider_factory: Some(build_acp_provider_factory(
app.config(),
app.secret_registry(),
)),
available_models: std::sync::Arc::new(parking_lot::RwLock::new(
if app.config().acp.available_models.is_empty() {
discover_models_from_config(app.config()).await
} else {
app.config().acp.available_models.clone()
},
)),
provider_names: acp_provider_names(app.config()),
mcp_manager: Some(Arc::clone(&mcp_manager_for_acp)),
auth_clients,
discovery_enabled: app.config().acp.discovery_enabled,
terminal_timeout_secs: app.config().acp.timeouts.terminal_secs,
project_rules: collect_project_rules(&app.skill_paths_for_registry()),
title_max_chars: app.config().memory.sessions.title_max_chars,
max_history: app.config().memory.sessions.max_history,
sqlite_path: Some(crate::db_url::resolve_db_url(app.config()).to_owned()),
session_data_dir: app
.config()
.session
.enabled
.then(|| std::path::PathBuf::from(&app.config().session.data_dir)),
ready_notification: None,
additional_directories: app.config().acp.additional_directories.clone(),
auth_methods: app.config().acp.auth_methods.clone(),
message_ids_enabled: app.config().acp.message_ids_enabled,
timeouts: app.config().acp.timeouts.clone(),
model_config: app.config().acp.model_config.clone(),
};
let shared_deps: Arc<RwLock<Option<Arc<SharedAgentDeps>>>> = Arc::new(RwLock::new(None));
let shared_deps_for_spawner = Arc::clone(&shared_deps);
let spawner: zeph_acp::SendAgentSpawner = Arc::new(move |channel, acp_ctx, session_ctx| {
let shared_deps = Arc::clone(&shared_deps_for_spawner);
Box::pin(async move {
let maybe_shared = shared_deps.read().await.clone();
let Some(shared) = maybe_shared else {
tracing::warn!("ACP request received before runtime became ready");
return;
};
Box::pin(spawn_acp_agent(shared, channel, acp_ctx, session_ctx)).await;
})
});
let mut state = zeph_acp::AcpHttpState::new(spawner, server_config);
match zeph_memory::store::SqliteStore::new(crate::db_url::resolve_db_url(app.config())).await {
Ok(store) => state = state.with_store(store),
Err(e) => tracing::warn!(error = %e, "failed to open SQLite for HTTP session endpoints"),
}
let router = zeph_acp::acp_router(state.clone());
let listener = tokio::net::TcpListener::bind(&bind_addr).await?;
tracing::info!("ACP HTTP server listening on {bind_addr}");
let server_task = tokio::spawn(async move { ::axum::serve(listener, router).await });
let (deps, _keepalive) =
match Box::pin(build_acp_deps(&app, None, Some(mcp_manager_for_acp))).await {
Ok(result) => result,
Err(err) => {
server_task.abort();
return Err(err);
}
};
let available_models = std::sync::Arc::clone(&deps.acp_available_models);
let provider = deps.provider.clone();
zeph_acp::warm_model_caches(provider, available_models).await;
*shared_deps.write().await = Some(Arc::new(deps));
state.mark_ready();
state.start_reaper();
tracing::info!("ACP server ready");
server_task.await??;
Ok(())
}
#[cfg(all(feature = "acp-http", feature = "session"))]
pub(crate) async fn build_combined_deps(
app: &AppBuilder,
supervisor: &std::sync::Arc<zeph_common::TaskSupervisor>,
) -> anyhow::Result<(
crate::serve::deps::ServeAgentDeps,
SharedAgentDeps,
Box<dyn std::any::Any>,
)> {
let core = build_shared_core(app, supervisor).await?;
let serve_deps = crate::serve::deps::assemble_serve_deps(app, &core, supervisor).await?;
let prebuilt_core = PrebuiltAcpCore {
core,
supervisor: std::sync::Arc::clone(supervisor),
};
let (acp_deps, keepalive) = Box::pin(build_acp_deps(app, Some(prebuilt_core), None)).await?;
Ok((serve_deps, acp_deps, keepalive))
}
#[cfg(all(feature = "acp-http", feature = "session"))]
pub(crate) fn acp_http_server_config(deps: &mut SharedAgentDeps) -> zeph_acp::AcpServerConfig {
zeph_acp::AcpServerConfig {
agent_name: deps.acp_agent_name.clone(),
agent_version: deps.acp_agent_version.clone(),
max_sessions: deps.acp_max_sessions,
session_idle_timeout_secs: deps.acp_session_idle_timeout_secs,
permission_file: deps.acp_permission_file.clone(),
provider_factory: deps.acp_provider_factory.take(),
available_models: std::sync::Arc::clone(&deps.acp_available_models),
provider_names: deps.acp_provider_names.clone(),
mcp_manager: Some(std::sync::Arc::clone(&deps.mcp_manager)),
auth_clients: deps.acp_auth_clients.clone(),
discovery_enabled: deps.acp_discovery_enabled,
terminal_timeout_secs: deps.acp_timeouts.terminal_secs,
project_rules: deps.acp_project_rules.clone(),
title_max_chars: deps.acp_title_max_chars,
max_history: deps.acp_max_history,
sqlite_path: Some(deps.sqlite_path.clone()),
session_data_dir: deps
.session_persistence_config
.enabled
.then(|| std::path::PathBuf::from(&deps.session_persistence_config.data_dir)),
ready_notification: None,
additional_directories: deps.acp_additional_directories.clone(),
auth_methods: deps.acp_auth_methods.clone(),
message_ids_enabled: deps.acp_message_ids_enabled,
timeouts: deps.acp_timeouts.clone(),
model_config: deps.acp_model_config.clone(),
}
}
#[cfg(all(feature = "acp-http", feature = "session"))]
pub(crate) async fn acp_http_ready_spawner(
deps: std::sync::Arc<SharedAgentDeps>,
) -> zeph_acp::SendAgentSpawner {
let available_models = std::sync::Arc::clone(&deps.acp_available_models);
let provider = deps.provider.clone();
zeph_acp::warm_model_caches(provider, available_models).await;
std::sync::Arc::new(move |channel, acp_ctx, session_ctx| {
let shared = std::sync::Arc::clone(&deps);
Box::pin(spawn_acp_agent(shared, channel, acp_ctx, session_ctx))
})
}
#[cfg(feature = "acp")]
pub(crate) fn print_acp_manifest() {
let manifest = serde_json::json!({
"name": env!("CARGO_PKG_NAME"),
"version": env!("CARGO_PKG_VERSION"),
"transport": "stdio",
"command": [env!("CARGO_PKG_NAME"), "--acp"],
"capabilities": ["prompt", "cancel", "load_session", "set_session_mode", "config_options", "ext_methods"],
"description": "Zeph AI Agent",
"readiness": {
"notification": {
"method": "zeph/ready",
"params": {
"version": env!("CARGO_PKG_VERSION"),
"pid": "<process-id>",
"log_file": "<configured-log-file>"
}
},
"http": {
"health_endpoint": "/health",
"statuses": [200, 503]
}
}
});
println!(
"{}",
serde_json::to_string_pretty(&manifest).unwrap_or_default()
);
}
#[cfg(all(test, feature = "acp"))]
mod tests {
use super::*;
use serial_test::serial;
use std::fs;
use std::sync::Arc;
use tempfile::TempDir;
use zeph_tools::executor::ToolExecutor;
#[derive(Default)]
struct TestVault {
secrets: std::collections::HashMap<String, String>,
erroring_keys: std::collections::HashSet<String>,
}
impl TestVault {
fn with_secret(mut self, key: &str, value: &str) -> Self {
self.secrets.insert(key.to_owned(), value.to_owned());
self
}
fn with_erroring_key(mut self, key: &str) -> Self {
self.erroring_keys.insert(key.to_owned());
self
}
}
impl zeph_core::vault::VaultProvider for TestVault {
fn get_secret(
&self,
key: &str,
) -> std::pin::Pin<
Box<
dyn std::future::Future<
Output = Result<Option<String>, zeph_core::vault::VaultError>,
> + Send
+ '_,
>,
> {
let result = if self.erroring_keys.contains(key) {
Err(zeph_core::vault::VaultError::Backend(
"simulated backend failure".to_owned(),
))
} else {
Ok(self.secrets.get(key).cloned())
};
Box::pin(async move { result })
}
}
fn acp_config_with(
auth_token: Option<&str>,
auth_clients: Vec<zeph_config::AcpAuthClient>,
) -> zeph_config::AcpConfig {
zeph_config::AcpConfig {
auth_token: auth_token.map(str::to_owned),
auth_clients,
..zeph_config::AcpConfig::default()
}
}
fn inline_client(id: &str, token: &str) -> zeph_config::AcpAuthClient {
zeph_config::AcpAuthClient {
id: id.to_owned(),
token: Some(token.to_owned()),
token_vault_key: None,
}
}
fn vault_client(id: &str, vault_key: &str) -> zeph_config::AcpAuthClient {
zeph_config::AcpAuthClient {
id: id.to_owned(),
token: None,
token_vault_key: Some(vault_key.to_owned()),
}
}
#[tokio::test]
async fn resolve_acp_auth_clients_empty_config_returns_empty() {
let cfg = acp_config_with(None, vec![]);
let clients = resolve_acp_auth_clients(&cfg, &TestVault::default())
.await
.unwrap();
assert!(clients.is_empty());
}
#[tokio::test]
async fn resolve_acp_auth_clients_legacy_token_becomes_default_client() {
let cfg = acp_config_with(Some("legacy-secret"), vec![]);
let clients = resolve_acp_auth_clients(&cfg, &TestVault::default())
.await
.unwrap();
assert_eq!(clients.len(), 1);
assert_eq!(clients[0].id, zeph_config::ACP_AUTH_CLIENT_ID_DEFAULT);
assert_eq!(clients[0].token, "legacy-secret");
}
#[tokio::test]
async fn resolve_acp_auth_clients_inline_token_resolved_directly() {
let cfg = acp_config_with(None, vec![inline_client("alice", "token-a")]);
let clients = resolve_acp_auth_clients(&cfg, &TestVault::default())
.await
.unwrap();
assert_eq!(clients.len(), 1);
assert_eq!(clients[0].id, "alice");
assert_eq!(clients[0].token, "token-a");
}
#[tokio::test]
async fn resolve_acp_auth_clients_vault_key_resolved_from_vault() {
let cfg = acp_config_with(None, vec![vault_client("alice", "ZEPH_ACP_TOKEN_ALICE")]);
let vault = TestVault::default().with_secret("ZEPH_ACP_TOKEN_ALICE", "vault-token-a");
let clients = resolve_acp_auth_clients(&cfg, &vault).await.unwrap();
assert_eq!(clients.len(), 1);
assert_eq!(clients[0].id, "alice");
assert_eq!(clients[0].token, "vault-token-a");
}
#[tokio::test]
async fn resolve_acp_auth_clients_missing_vault_key_soft_disables_client() {
let cfg = acp_config_with(
None,
vec![
vault_client("alice", "ZEPH_ACP_TOKEN_ALICE"),
inline_client("bob", "token-b"),
],
);
let clients = resolve_acp_auth_clients(&cfg, &TestVault::default())
.await
.unwrap();
assert_eq!(clients.len(), 1);
assert_eq!(clients[0].id, "bob");
}
#[tokio::test]
async fn resolve_acp_auth_clients_vault_backend_error_soft_disables_client() {
let cfg = acp_config_with(
None,
vec![
vault_client("alice", "ZEPH_ACP_TOKEN_ALICE"),
inline_client("bob", "token-b"),
],
);
let vault = TestVault::default().with_erroring_key("ZEPH_ACP_TOKEN_ALICE");
let clients = resolve_acp_auth_clients(&cfg, &vault).await.unwrap();
assert_eq!(clients.len(), 1);
assert_eq!(clients[0].id, "bob");
}
#[tokio::test]
async fn resolve_acp_auth_clients_empty_vault_token_soft_disables_client() {
let cfg = acp_config_with(
None,
vec![
vault_client("alice", "ZEPH_ACP_TOKEN_ALICE"),
inline_client("bob", "token-b"),
],
);
let vault = TestVault::default().with_secret("ZEPH_ACP_TOKEN_ALICE", "");
let clients = resolve_acp_auth_clients(&cfg, &vault).await.unwrap();
assert_eq!(clients.len(), 1);
assert_eq!(clients[0].id, "bob");
}
#[tokio::test]
async fn resolve_acp_auth_clients_sole_empty_vault_token_fails_closed() {
let cfg = acp_config_with(None, vec![vault_client("alice", "ZEPH_ACP_TOKEN_ALICE")]);
let vault = TestVault::default().with_secret("ZEPH_ACP_TOKEN_ALICE", "");
let err = resolve_acp_auth_clients(&cfg, &vault).await.unwrap_err();
assert!(
err.to_string().contains("refusing to start"),
"unexpected error message: {err}"
);
}
#[tokio::test]
async fn resolve_acp_auth_clients_sole_missing_vault_key_fails_closed() {
let cfg = acp_config_with(None, vec![vault_client("alice", "ZEPH_ACP_TOKEN_ALICE")]);
let err = resolve_acp_auth_clients(&cfg, &TestVault::default())
.await
.unwrap_err();
assert!(
err.to_string().contains("refusing to start"),
"unexpected error message: {err}"
);
}
#[tokio::test]
async fn resolve_acp_auth_clients_sole_client_backend_error_fails_closed() {
let cfg = acp_config_with(None, vec![vault_client("alice", "ZEPH_ACP_TOKEN_ALICE")]);
let vault = TestVault::default().with_erroring_key("ZEPH_ACP_TOKEN_ALICE");
let err = resolve_acp_auth_clients(&cfg, &vault).await.unwrap_err();
assert!(
err.to_string().contains("refusing to start"),
"unexpected error message: {err}"
);
}
#[tokio::test]
async fn resolve_acp_auth_clients_legacy_token_alone_never_empties_so_no_fail_closed_path() {
let cfg = acp_config_with(Some("legacy-secret"), vec![]);
let clients = resolve_acp_auth_clients(&cfg, &TestVault::default())
.await
.unwrap();
assert_eq!(clients.len(), 1);
}
#[tokio::test]
async fn resolve_acp_auth_clients_rejects_vault_token_colliding_with_inline_token() {
let cfg = acp_config_with(
None,
vec![
inline_client("alice", "shared-secret"),
vault_client("bob", "ZEPH_ACP_TOKEN_BOB"),
],
);
let vault = TestVault::default().with_secret("ZEPH_ACP_TOKEN_BOB", "shared-secret");
let err = resolve_acp_auth_clients(&cfg, &vault).await.unwrap_err();
assert!(
err.to_string().contains("collides"),
"unexpected error: {err}"
);
}
#[tokio::test]
async fn resolve_acp_auth_clients_rejects_two_vault_tokens_resolving_to_same_secret() {
let cfg = acp_config_with(
None,
vec![
vault_client("alice", "ZEPH_ACP_TOKEN_ALICE"),
vault_client("bob", "ZEPH_ACP_TOKEN_BOB"),
],
);
let vault = TestVault::default()
.with_secret("ZEPH_ACP_TOKEN_ALICE", "same-secret")
.with_secret("ZEPH_ACP_TOKEN_BOB", "same-secret");
let err = resolve_acp_auth_clients(&cfg, &vault).await.unwrap_err();
assert!(
err.to_string().contains("collides"),
"unexpected error: {err}"
);
}
#[tokio::test]
async fn resolve_acp_auth_clients_rejects_vault_token_colliding_with_legacy_default() {
let cfg = acp_config_with(
Some("legacy-secret"),
vec![vault_client("alice", "ZEPH_ACP_TOKEN_ALICE")],
);
let vault = TestVault::default().with_secret("ZEPH_ACP_TOKEN_ALICE", "legacy-secret");
let err = resolve_acp_auth_clients(&cfg, &vault).await.unwrap_err();
assert!(
err.to_string().contains("collides"),
"unexpected error: {err}"
);
}
fn make_rules_dir(dir: &std::path::Path, files: &[&str]) {
let rules = dir.join(".claude").join("rules");
fs::create_dir_all(&rules).unwrap();
for name in files {
fs::write(rules.join(name), b"").unwrap();
}
}
#[test]
#[serial]
fn collect_project_rules_empty_skill_paths_no_rules_dir() {
let tmp = TempDir::new().unwrap();
let orig = std::env::current_dir().unwrap();
std::env::set_current_dir(tmp.path()).unwrap();
let result = collect_project_rules(&[]);
std::env::set_current_dir(orig).unwrap();
assert!(result.is_empty());
}
#[test]
#[serial]
fn collect_project_rules_picks_md_files_from_rules_dir() {
let tmp = TempDir::new().unwrap();
make_rules_dir(tmp.path(), &["rust-code.md", "testing.md", "notes.txt"]);
let orig = std::env::current_dir().unwrap();
std::env::set_current_dir(tmp.path()).unwrap();
let result = collect_project_rules(&[]);
std::env::set_current_dir(orig).unwrap();
assert_eq!(result.len(), 2);
let names: Vec<_> = result
.iter()
.filter_map(|p| p.file_name())
.map(|n| n.to_string_lossy().into_owned())
.collect();
assert!(names.contains(&"rust-code.md".to_owned()));
assert!(names.contains(&"testing.md".to_owned()));
assert!(!names.contains(&"notes.txt".to_owned()));
}
#[test]
#[serial]
fn collect_project_rules_includes_skill_files() {
let tmp = TempDir::new().unwrap();
let skill_file = tmp.path().join("my-skill.md");
fs::write(&skill_file, b"").unwrap();
let skill_dir = tmp.path().join("skills-dir");
fs::create_dir_all(&skill_dir).unwrap();
let orig = std::env::current_dir().unwrap();
std::env::set_current_dir(tmp.path()).unwrap();
let result = collect_project_rules(&[skill_file.clone(), skill_dir]);
std::env::set_current_dir(orig).unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0], skill_file);
}
#[tokio::test]
async fn diagnostics_tool_call_dispatches_through_acp_composite_chain() {
let config = zeph_core::config::Config::default();
let file_executor = zeph_tools::FileExecutor::new(vec![]);
let shell_executor = zeph_tools::ShellExecutor::new(&config.tools.shell);
let scrape_executor = zeph_tools::WebScrapeExecutor::new(&config.tools.scrape);
let diagnostics_executor = crate::agent_setup::build_diagnostics_executor(&config);
let base_executor = crate::agent_setup::build_base_executor_chain(
file_executor,
shell_executor,
scrape_executor,
diagnostics_executor,
zeph_tools::GetCurrentTimeExecutor::default(),
vec![],
);
let policy =
zeph_tools::PermissionPolicy::default().with_autonomy(zeph_tools::AutonomyLevel::Full);
let base_executor = zeph_tools::TrustGateExecutor::new(base_executor, policy);
let outside = std::env::temp_dir();
let mut params = serde_json::Map::new();
params.insert(
"path".into(),
serde_json::Value::String(outside.display().to_string()),
);
let call = zeph_tools::ToolCall {
tool_id: "diagnostics".into(),
params,
caller_id: None,
context: None,
tool_call_id: String::new(),
skill_name: None,
};
let result = base_executor.execute_tool_call(&call).await;
assert!(
matches!(result, Err(zeph_tools::ToolError::SandboxViolation { .. })),
"expected SandboxViolation from DiagnosticsExecutor, got {result:?}"
);
}
#[tokio::test]
async fn diagnostics_requires_confirmation_in_acp_composite_chain() {
let config = zeph_core::config::Config::default();
let file_executor = zeph_tools::FileExecutor::new(vec![]);
let shell_executor = zeph_tools::ShellExecutor::new(&config.tools.shell);
let scrape_executor = zeph_tools::WebScrapeExecutor::new(&config.tools.scrape);
let diagnostics_executor = crate::agent_setup::build_diagnostics_executor(&config);
let base_executor = crate::agent_setup::build_base_executor_chain(
file_executor,
shell_executor,
scrape_executor,
diagnostics_executor,
zeph_tools::GetCurrentTimeExecutor::default(),
vec![],
);
let base_executor = zeph_tools::TrustGateExecutor::new(
base_executor,
zeph_tools::PermissionPolicy::default(),
);
let call = zeph_tools::ToolCall {
tool_id: "diagnostics".into(),
params: serde_json::Map::new(),
caller_id: None,
context: None,
tool_call_id: String::new(),
skill_name: None,
};
let result = base_executor.execute_tool_call(&call).await;
assert!(
matches!(
result,
Err(zeph_tools::ToolError::ConfirmationRequired { .. })
),
"expected ConfirmationRequired for diagnostics under Supervised autonomy, got {result:?}"
);
}
#[derive(Debug)]
struct AcpTaggedMock(&'static str);
impl zeph_tools::executor::ToolExecutor for AcpTaggedMock {
async fn execute(
&self,
_response: &str,
) -> Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError> {
Ok(None)
}
async fn execute_tool_call(
&self,
call: &zeph_tools::ToolCall,
) -> Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError> {
if call.tool_id != self.0 {
return Ok(None);
}
Ok(Some(zeph_tools::ToolOutput {
tool_name: call.tool_id.clone(),
summary: "ok".into(),
blocks_executed: 1,
filter_stats: None,
diff: None,
streamed: false,
terminal_id: None,
locations: None,
raw_response: None,
claim_source: None,
..Default::default()
}))
}
zeph_tools::tool_executor_no_inner_defaults!();
}
fn acp_test_call(tool_id: &str) -> zeph_tools::ToolCall {
zeph_tools::ToolCall {
tool_id: tool_id.into(),
params: serde_json::Map::new(),
caller_id: None,
context: None,
tool_call_id: String::new(),
skill_name: None,
}
}
#[tokio::test]
async fn quarantine_blocks_memory_and_mcp_in_acp_composite_chain() {
let mcp_tool = zeph_mcp::McpTool {
server_id: "mcp".to_owned(),
name: "write_file".to_owned(),
description: String::new(),
input_schema: serde_json::Value::Null,
output_schema: None,
security_meta: zeph_mcp::tool::ToolSecurityMeta::default(),
};
let mcp_tool_id = mcp_tool.sanitized_id();
assert_eq!(mcp_tool_id, "mcp_write_file");
let base_tool = zeph_tools::CompositeExecutor::new(
AcpTaggedMock("read"),
AcpTaggedMock("mcp_write_file"),
);
let inner_executor =
zeph_tools::DynExecutor(std::sync::Arc::new(zeph_tools::CompositeExecutor::new(
AcpTaggedMock("load_skill"),
zeph_tools::CompositeExecutor::new(
AcpTaggedMock("memory_save"),
zeph_tools::CompositeExecutor::new(AcpTaggedMock("overflow_flush"), base_tool),
),
)));
let (gated, mcp_ids_handle) = crate::agent_setup::apply_common_tool_gating(
inner_executor,
&zeph_tools::PermissionPolicy::default(),
);
crate::agent_setup::register_mcp_tool_ids(&mcp_ids_handle, std::slice::from_ref(&mcp_tool));
zeph_tools::executor::ToolExecutor::set_effective_trust(
&gated,
zeph_common::SkillTrustLevel::Quarantined,
);
let memory_result = gated.execute_tool_call(&acp_test_call("memory_save")).await;
assert!(
matches!(memory_result, Err(zeph_tools::ToolError::Blocked { .. })),
"memory_save must be denied under Quarantine, got {memory_result:?}"
);
let mcp_result = gated.execute_tool_call(&acp_test_call(&mcp_tool_id)).await;
assert!(
matches!(mcp_result, Err(zeph_tools::ToolError::Blocked { .. })),
"MCP-sourced tool must be denied under Quarantine, got {mcp_result:?}"
);
let skill_load_result = gated.execute_tool_call(&acp_test_call("load_skill")).await;
assert!(
matches!(
skill_load_result,
Err(zeph_tools::ToolError::Blocked { .. })
),
"load_skill must be denied under Quarantine, got {skill_load_result:?}"
);
let read_result = gated.execute_tool_call(&acp_test_call("read")).await;
assert!(
read_result.is_ok(),
"readonly native tool must remain reachable under Quarantine, got {read_result:?}"
);
}
#[tokio::test]
async fn policy_gate_denies_tool_in_acp_composite_chain() {
let config = zeph_core::config::Config::default();
let file_executor = zeph_tools::FileExecutor::new(vec![]);
let shell_executor = zeph_tools::ShellExecutor::new(&config.tools.shell);
let scrape_executor = zeph_tools::WebScrapeExecutor::new(&config.tools.scrape);
let diagnostics_executor = crate::agent_setup::build_diagnostics_executor(&config);
let base_executor = crate::agent_setup::build_base_executor_chain(
file_executor,
shell_executor,
scrape_executor,
diagnostics_executor,
zeph_tools::GetCurrentTimeExecutor::default(),
vec![],
);
let policy =
zeph_tools::PermissionPolicy::default().with_autonomy(zeph_tools::AutonomyLevel::Full);
let base_executor = zeph_tools::TrustGateExecutor::new(base_executor, policy);
let policy_config = zeph_tools::PolicyConfig {
enabled: true,
default_effect: zeph_tools::DefaultEffect::Allow,
rules: vec![zeph_tools::PolicyRuleConfig {
effect: zeph_tools::PolicyEffect::Deny,
tool: "diagnostics".into(),
paths: vec![],
env: vec![],
trust_level: None,
args_match: None,
capabilities: vec![],
}],
..Default::default()
};
let enforcer = zeph_tools::PolicyEnforcer::compile(&policy_config).unwrap();
let policy_context = std::sync::Arc::new(RwLock::new(zeph_tools::PolicyContext {
trust_level: zeph_common::SkillTrustLevel::Trusted,
env: std::collections::HashMap::new(),
}));
let gated = zeph_tools::PolicyGateExecutor::new(
base_executor,
std::sync::Arc::new(enforcer),
policy_context,
);
let call = zeph_tools::ToolCall {
tool_id: "diagnostics".into(),
params: serde_json::Map::new(),
caller_id: None,
context: None,
tool_call_id: String::new(),
skill_name: None,
};
let result = gated.execute_tool_call(&call).await;
assert!(
matches!(result, Err(zeph_tools::ToolError::Blocked { .. })),
"expected Blocked from PolicyGateExecutor deny rule, got {result:?}"
);
}
#[tokio::test]
async fn adversarial_policy_gate_denies_tool_in_acp_composite_chain() {
struct AlwaysDenyLlm;
impl zeph_tools::PolicyLlmClient for AlwaysDenyLlm {
fn chat<'a>(
&'a self,
_messages: &'a [zeph_tools::PolicyMessage],
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<String, String>> + Send + 'a>,
> {
Box::pin(async move { Ok("DENY: test policy".to_owned()) })
}
}
let config = zeph_core::config::Config::default();
let file_executor = zeph_tools::FileExecutor::new(vec![]);
let shell_executor = zeph_tools::ShellExecutor::new(&config.tools.shell);
let scrape_executor = zeph_tools::WebScrapeExecutor::new(&config.tools.scrape);
let diagnostics_executor = crate::agent_setup::build_diagnostics_executor(&config);
let base_executor = crate::agent_setup::build_base_executor_chain(
file_executor,
shell_executor,
scrape_executor,
diagnostics_executor,
zeph_tools::GetCurrentTimeExecutor::default(),
vec![],
);
let policy =
zeph_tools::PermissionPolicy::default().with_autonomy(zeph_tools::AutonomyLevel::Full);
let base_executor = zeph_tools::TrustGateExecutor::new(base_executor, policy);
let validator = std::sync::Arc::new(zeph_tools::PolicyValidator::new(
vec!["never allow diagnostics".to_owned()],
std::time::Duration::from_millis(500),
false,
vec![],
));
let llm_client: std::sync::Arc<dyn zeph_tools::PolicyLlmClient> =
std::sync::Arc::new(AlwaysDenyLlm);
let gated =
zeph_tools::AdversarialPolicyGateExecutor::new(base_executor, validator, llm_client);
let call = zeph_tools::ToolCall {
tool_id: "diagnostics".into(),
params: serde_json::Map::new(),
caller_id: None,
context: None,
tool_call_id: String::new(),
skill_name: None,
};
let result = gated.execute_tool_call(&call).await;
assert!(
matches!(result, Err(zeph_tools::ToolError::Blocked { .. })),
"expected Blocked from AdversarialPolicyGateExecutor deny decision, got {result:?}"
);
}
#[tokio::test]
async fn policy_and_quarantine_trust_gate_both_enforce_in_acp_composite_chain() {
use zeph_tools::executor::ToolExecutor;
let mcp_tool = zeph_mcp::McpTool {
server_id: "mcp".to_owned(),
name: "write_file".to_owned(),
description: String::new(),
input_schema: serde_json::Value::Null,
output_schema: None,
security_meta: zeph_mcp::tool::ToolSecurityMeta::default(),
};
let base_tool = zeph_tools::CompositeExecutor::new(
AcpTaggedMock("read"),
AcpTaggedMock("mcp_write_file"),
);
let inner_executor =
zeph_tools::DynExecutor(std::sync::Arc::new(zeph_tools::CompositeExecutor::new(
AcpTaggedMock("load_skill"),
zeph_tools::CompositeExecutor::new(
AcpTaggedMock("memory_save"),
zeph_tools::CompositeExecutor::new(AcpTaggedMock("overflow_flush"), base_tool),
),
)));
let (trust_gated, mcp_ids_handle) = crate::agent_setup::apply_common_tool_gating(
inner_executor,
&zeph_tools::PermissionPolicy::default(),
);
crate::agent_setup::register_mcp_tool_ids(&mcp_ids_handle, std::slice::from_ref(&mcp_tool));
zeph_tools::ToolExecutor::set_effective_trust(
&trust_gated,
zeph_common::SkillTrustLevel::Quarantined,
);
let policy_config = zeph_tools::PolicyConfig {
enabled: true,
default_effect: zeph_tools::DefaultEffect::Allow,
rules: vec![zeph_tools::PolicyRuleConfig {
effect: zeph_tools::PolicyEffect::Deny,
tool: "overflow_flush".into(),
paths: vec![],
env: vec![],
trust_level: None,
args_match: None,
capabilities: vec![],
}],
..Default::default()
};
let enforcer = zeph_tools::PolicyEnforcer::compile(&policy_config).unwrap();
let policy_context = std::sync::Arc::new(RwLock::new(zeph_tools::PolicyContext {
trust_level: zeph_common::SkillTrustLevel::Trusted,
env: std::collections::HashMap::new(),
}));
let gated = zeph_tools::PolicyGateExecutor::new(
trust_gated,
std::sync::Arc::new(enforcer),
policy_context,
);
let policy_denied = gated
.execute_tool_call(&acp_test_call("overflow_flush"))
.await;
assert!(
matches!(policy_denied, Err(zeph_tools::ToolError::Blocked { .. })),
"expected Blocked from PolicyGateExecutor's own deny rule, got {policy_denied:?}"
);
let quarantine_denied = gated.execute_tool_call(&acp_test_call("load_skill")).await;
assert!(
matches!(
quarantine_denied,
Err(zeph_tools::ToolError::Blocked { .. })
),
"expected Blocked from TrustGateExecutor's Quarantine enforcement, got {quarantine_denied:?}"
);
let allowed = gated.execute_tool_call(&acp_test_call("read")).await;
assert!(
allowed.is_ok(),
"expected read to dispatch normally through the merged gate stack, got {allowed:?}"
);
}
#[tokio::test]
async fn capability_scopes_denies_tool_outside_scope_in_acp_composite_chain() {
use std::collections::HashSet;
use zeph_tools::scope::build_scoped_executor;
let config = zeph_core::config::Config::default();
let file_executor = zeph_tools::FileExecutor::new(vec![]);
let shell_executor = zeph_tools::ShellExecutor::new(&config.tools.shell);
let scrape_executor = zeph_tools::WebScrapeExecutor::new(&config.tools.scrape);
let diagnostics_executor = crate::agent_setup::build_diagnostics_executor(&config);
let base_executor = crate::agent_setup::build_base_executor_chain(
file_executor,
shell_executor,
scrape_executor,
diagnostics_executor,
zeph_tools::GetCurrentTimeExecutor::default(),
vec![],
);
let policy =
zeph_tools::PermissionPolicy::default().with_autonomy(zeph_tools::AutonomyLevel::Full);
let base_executor = zeph_tools::TrustGateExecutor::new(base_executor, policy);
let registry_ids: HashSet<String> = base_executor
.tool_definitions()
.into_iter()
.map(|def| {
let id = def.id.to_string();
if id.contains(':') {
id
} else {
format!("builtin:{id}")
}
})
.collect();
let scopes_cfg = zeph_config::CapabilityScopesConfig {
default_scope: "narrow".to_owned(),
scopes: std::collections::HashMap::from([(
"narrow".to_owned(),
zeph_config::ScopeConfig {
patterns: vec!["builtin:read".to_owned()],
},
)]),
..Default::default()
};
let scoped = build_scoped_executor(base_executor, &scopes_cfg, ®istry_ids)
.expect("build_scoped_executor must compile a valid single-pattern scope");
let denied = scoped
.execute_tool_call(&acp_test_call("diagnostics"))
.await;
assert!(
matches!(denied, Err(zeph_tools::ToolError::OutOfScope { .. })),
"expected OutOfScope from ScopedToolExecutor for a tool outside the configured \
scope, got {denied:?}"
);
let allowed = scoped.execute_tool_call(&acp_test_call("read")).await;
assert!(
!matches!(allowed, Err(zeph_tools::ToolError::OutOfScope { .. })),
"expected read to reach past ScopedToolExecutor since it matches the active \
scope's pattern, got {allowed:?}"
);
}
#[tokio::test]
async fn shadow_probe_executor_reaches_shadow_sentinel_in_acp_composite_chain() {
use zeph_core::agent::shadow_sentinel::{
ProbeVerdict, SafetyProbe, SentinelEvent, ShadowEventStore, ShadowSentinel,
};
use zeph_tools::{ProbeGate, ToolCall, ToolOutput};
struct AllowProbe;
impl SafetyProbe for AllowProbe {
fn evaluate<'a>(
&'a self,
_: &'a str,
_: &'a serde_json::Value,
_: &'a [SentinelEvent],
) -> std::pin::Pin<Box<dyn std::future::Future<Output = ProbeVerdict> + Send + 'a>>
{
Box::pin(async { ProbeVerdict::Allow })
}
}
struct OkExec;
impl ToolExecutor for OkExec {
async fn execute(&self, _: &str) -> Result<Option<ToolOutput>, zeph_tools::ToolError> {
Ok(None)
}
async fn execute_tool_call(
&self,
call: &ToolCall,
) -> Result<Option<ToolOutput>, zeph_tools::ToolError> {
Ok(Some(ToolOutput {
tool_name: call.tool_id.clone(),
summary: "command completed".to_owned(),
blocks_executed: 1,
filter_stats: None,
diff: None,
streamed: false,
terminal_id: None,
locations: None,
raw_response: None,
claim_source: None,
..Default::default()
}))
}
zeph_tools::tool_executor_no_inner_defaults!();
}
let pool = zeph_db::DbConfig {
url: ":memory:".to_owned(),
..Default::default()
}
.connect()
.await
.expect("connect + migrate in-memory sqlite pool");
let sentinel = std::sync::Arc::new(ShadowSentinel::new(
ShadowEventStore::new(pool.clone()),
Box::new(AllowProbe),
zeph_config::ShadowSentinelConfig {
enabled: true,
..Default::default()
},
"acp-conversation-42",
));
let probe_gate: std::sync::Arc<dyn ProbeGate> =
std::sync::Arc::new(crate::runner::ShadowSentinelProbeGateAdapter {
sentinel: std::sync::Arc::clone(&sentinel),
});
let executor = zeph_tools::ShadowProbeExecutor::new(
OkExec,
probe_gate,
std::sync::Arc::new(std::sync::atomic::AtomicU64::new(1)),
std::sync::Arc::new(parking_lot::RwLock::new("calm".to_owned())),
);
let result = executor
.execute_tool_call(&acp_test_call("builtin:shell"))
.await;
assert!(result.unwrap().is_some(), "tool call must succeed");
sentinel.drain_pending().await;
let events = ShadowEventStore::new(pool)
.get_trajectory("acp-conversation-42", 10)
.await
.expect("get_trajectory must succeed against the in-memory pool");
assert!(
events.iter().any(|e| e.event_type == "tool_call"
&& e.context_summary.as_deref() == Some("command completed")),
"expected the ShadowProbeExecutor-driven tool_call event to be persisted via ACP's \
reused ShadowSentinelProbeGateAdapter, got: {events:?}"
);
}
#[tokio::test]
async fn trajectory_signal_queue_receives_policy_denial_in_acp_composite_chain() {
let config = zeph_core::config::Config::default();
let file_executor = zeph_tools::FileExecutor::new(vec![]);
let shell_executor = zeph_tools::ShellExecutor::new(&config.tools.shell);
let scrape_executor = zeph_tools::WebScrapeExecutor::new(&config.tools.scrape);
let diagnostics_executor = crate::agent_setup::build_diagnostics_executor(&config);
let base_executor = crate::agent_setup::build_base_executor_chain(
file_executor,
shell_executor,
scrape_executor,
diagnostics_executor,
zeph_tools::GetCurrentTimeExecutor::default(),
vec![],
);
let (trust_gated, _mcp_ids_handle) = crate::agent_setup::apply_common_tool_gating(
zeph_tools::DynExecutor(Arc::new(base_executor)),
&zeph_tools::PermissionPolicy::default().with_autonomy(zeph_tools::AutonomyLevel::Full),
);
let policy_config = zeph_tools::PolicyConfig {
enabled: true,
default_effect: zeph_tools::DefaultEffect::Allow,
rules: vec![zeph_tools::PolicyRuleConfig {
effect: zeph_tools::PolicyEffect::Deny,
tool: "overflow_flush".into(),
paths: vec![],
env: vec![],
trust_level: None,
args_match: None,
capabilities: vec![],
}],
..Default::default()
};
let enforcer = zeph_tools::PolicyEnforcer::compile(&policy_config).unwrap();
let pieces = crate::agent_setup::PolicyGatePieces {
policy_enforcer: Some(Arc::new(enforcer)),
adversarial_validator: None,
adversarial_llm_client: None,
adv_policy_info: None,
policy_configured: true,
};
let trajectory_risk_slot: zeph_tools::TrajectoryRiskSlot =
Arc::new(parking_lot::RwLock::new(0u8));
let trajectory_signal_queue: zeph_tools::RiskSignalQueue =
Arc::new(parking_lot::Mutex::new(Vec::new()));
let gated = crate::agent_setup::apply_policy_gate_chain(
trust_gated,
&pieces,
None,
Some((&trajectory_risk_slot, &trajectory_signal_queue)),
);
let denied = gated
.execute_tool_call(&acp_test_call("overflow_flush"))
.await;
assert!(
matches!(denied, Err(zeph_tools::ToolError::Blocked { .. })),
"expected Blocked from PolicyGateExecutor's deny rule, got {denied:?}"
);
assert_eq!(
*trajectory_signal_queue.lock(),
vec![1u8],
"expected the PolicyDeny signal code (1) to be pushed into the shared trajectory \
signal queue after a denied tool call, proving apply_policy_gate_chain's ACP call \
site actually wires PolicyGateExecutor::with_signal_queue instead of passing None"
);
}
#[tokio::test]
async fn trajectory_signal_queue_receives_scope_denial_in_acp_composite_chain() {
use std::collections::HashSet;
use zeph_tools::scope::build_scoped_executor;
let config = zeph_core::config::Config::default();
let file_executor = zeph_tools::FileExecutor::new(vec![]);
let shell_executor = zeph_tools::ShellExecutor::new(&config.tools.shell);
let scrape_executor = zeph_tools::WebScrapeExecutor::new(&config.tools.scrape);
let diagnostics_executor = crate::agent_setup::build_diagnostics_executor(&config);
let base_executor = crate::agent_setup::build_base_executor_chain(
file_executor,
shell_executor,
scrape_executor,
diagnostics_executor,
zeph_tools::GetCurrentTimeExecutor::default(),
vec![],
);
let policy =
zeph_tools::PermissionPolicy::default().with_autonomy(zeph_tools::AutonomyLevel::Full);
let base_executor = zeph_tools::TrustGateExecutor::new(base_executor, policy);
let registry_ids: HashSet<String> = base_executor
.tool_definitions()
.into_iter()
.map(|def| {
let id = def.id.to_string();
if id.contains(':') {
id
} else {
format!("builtin:{id}")
}
})
.collect();
let scopes_cfg = zeph_config::CapabilityScopesConfig {
default_scope: "narrow".to_owned(),
scopes: std::collections::HashMap::from([(
"narrow".to_owned(),
zeph_config::ScopeConfig {
patterns: vec!["builtin:read".to_owned()],
},
)]),
..Default::default()
};
let scoped = build_scoped_executor(base_executor, &scopes_cfg, ®istry_ids)
.expect("build_scoped_executor must compile a valid single-pattern scope");
let trajectory_signal_queue: zeph_tools::RiskSignalQueue =
Arc::new(parking_lot::Mutex::new(Vec::new()));
let scoped = scoped.with_signal_queue(Arc::clone(&trajectory_signal_queue));
let denied = scoped
.execute_tool_call(&acp_test_call("diagnostics"))
.await;
assert!(
matches!(denied, Err(zeph_tools::ToolError::OutOfScope { .. })),
"expected OutOfScope from ScopedToolExecutor for a tool outside the configured \
scope, got {denied:?}"
);
assert_eq!(
*trajectory_signal_queue.lock(),
vec![3u8],
"expected the OutOfScope signal code (3) to be pushed into the shared trajectory \
signal queue after a scope-denied tool call, proving spawn_acp_agent's new \
`.with_signal_queue(...)` call on the ScopedToolExecutor branch is reachable"
);
}
#[tokio::test]
async fn invoke_skill_reaches_skill_invoke_executor_in_full_acp_session_composite() {
let (session_composite, _trust_snapshot) =
build_full_acp_session_composite_with_native_fs_shell().await;
let mut params = serde_json::Map::new();
params.insert(
"skill_name".to_owned(),
serde_json::Value::String("nonexistent-skill".to_owned()),
);
let call = zeph_tools::ToolCall {
tool_id: "invoke_skill".into(),
params,
caller_id: None,
context: None,
tool_call_id: String::new(),
skill_name: None,
};
let result = session_composite.execute_tool_call_erased(&call).await;
let output = result
.expect("invoke_skill must dispatch successfully through SkillInvokeExecutor")
.expect("SkillInvokeExecutor must always return Some(ToolOutput) for invoke_skill");
assert!(
output
.summary
.contains("skill not found: nonexistent-skill"),
"expected the \"skill not found: ...\" summary that only SkillInvokeExecutor \
produces, proving invoke_skill actually reaches it in the full ACP session \
composite instead of falling through to memory/overflow/base, got: {output:?}"
);
}
#[tokio::test]
async fn acp_memory_maintenance_loops_registered_on_connection_supervisor() {
let mock_provider =
zeph_llm::any::AnyProvider::Mock(zeph_llm::mock::MockProvider::default());
let memory = std::sync::Arc::new(
zeph_memory::semantic::SemanticMemory::new(
":memory:",
"http://127.0.0.1:1",
None,
mock_provider.clone(),
"test",
)
.await
.unwrap(),
);
let mut config = zeph_core::config::Config::default();
config.memory.compression_guidelines.enabled = true;
config.memory.tree.enabled = true;
config.memory.hebbian.enabled = true;
config.memory.episodic_consolidation.enabled = true;
config.memory.optical_forgetting.enabled = true;
let app = crate::bootstrap::AppBuilder::for_test(config);
let cancel = tokio_util::sync::CancellationToken::new();
let supervisor = zeph_common::TaskSupervisor::new(cancel);
agent_setup::spawn_memory_maintenance_loops(
&app,
&memory,
&mock_provider,
&supervisor,
None,
false,
"acp",
);
let names: std::collections::HashSet<String> = supervisor
.snapshot()
.into_iter()
.map(|s| s.name.to_string())
.collect();
for expected in [
"mem-eviction",
"mem-tier-promotion",
"mem-scene-consolidation",
"mem-consolidation",
"mem-forgetting",
"mem-guidelines",
"mem-tree-consolidation",
"mem-hebbian-consolidation",
"mem-episodic-consolidation",
"mem-optical-forgetting",
] {
assert!(
names.contains(expected),
"expected {expected} registered on the ACP connection's memory supervisor, \
got {names:?}"
);
}
}
#[derive(Debug)]
struct AcpNativeStandIn {
tool_id: &'static str,
}
impl ToolExecutor for AcpNativeStandIn {
async fn execute(
&self,
_response: &str,
) -> Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError> {
Ok(None)
}
async fn execute_tool_call(
&self,
call: &zeph_tools::ToolCall,
) -> Result<Option<zeph_tools::ToolOutput>, zeph_tools::ToolError> {
if call.tool_id != self.tool_id {
return Ok(None);
}
panic!(
"AcpNativeStandIn({}) reached — gate did not intercept",
self.tool_id
);
}
zeph_tools::tool_executor_no_inner_defaults!();
}
async fn build_full_acp_session_composite_with_native_fs_shell() -> (
Arc<dyn ErasedToolExecutor>,
Arc<
RwLock<std::collections::HashMap<String, zeph_core::skill_invoker::SkillTrustSnapshot>>,
>,
) {
let registry = Arc::new(RwLock::new(zeph_skills::registry::SkillRegistry::empty()));
let (skill_loader_executor, skill_invoke_executor, trust_snapshot) =
agent_setup::build_skill_executors(®istry);
let mock_provider =
zeph_llm::any::AnyProvider::Mock(zeph_llm::mock::MockProvider::default());
let memory = Arc::new(
zeph_memory::semantic::SemanticMemory::new(
":memory:",
"http://127.0.0.1:1",
None,
mock_provider,
"test",
)
.await
.unwrap(),
);
let memory_executor = zeph_core::memory_tools::MemoryToolExecutor::with_validator(
Arc::clone(&memory),
zeph_memory::ConversationId(0),
zeph_sanitizer::memory_validation::MemoryWriteValidator::new(
zeph_core::config::Config::default()
.security
.memory_validation
.clone(),
),
);
let overflow_executor =
zeph_core::overflow_tools::OverflowToolExecutor::new(Arc::new(memory.sqlite().clone()));
let mut base: Arc<dyn ErasedToolExecutor> = Arc::new(zeph_tools::FileExecutor::new(vec![]));
let filtered =
zeph_tools::ToolFilter::new(zeph_tools::DynExecutor(base), &["read", "write", "glob"]);
base = Arc::new(zeph_tools::CompositeExecutor::new(
AcpNativeStandIn {
tool_id: "write_file",
},
filtered,
));
base = Arc::new(zeph_tools::CompositeExecutor::new(
AcpNativeStandIn { tool_id: "bash" },
zeph_tools::DynExecutor(base),
));
base = Arc::new(zeph_tools::CompositeExecutor::new(
skill_loader_executor,
zeph_tools::CompositeExecutor::new(
skill_invoke_executor,
zeph_tools::CompositeExecutor::new(
memory_executor,
zeph_tools::CompositeExecutor::new(
overflow_executor,
zeph_tools::DynExecutor(base),
),
),
),
));
(base, trust_snapshot)
}
#[tokio::test]
async fn policy_gate_denies_skill_and_memory_tools_in_full_acp_session_composite() {
let (session_composite, _trust_snapshot) =
build_full_acp_session_composite_with_native_fs_shell().await;
let policy_config = zeph_tools::PolicyConfig {
enabled: true,
default_effect: zeph_tools::DefaultEffect::Allow,
rules: ["load_skill", "memory_search", "write_file", "bash"]
.into_iter()
.map(|tool| zeph_tools::PolicyRuleConfig {
effect: zeph_tools::PolicyEffect::Deny,
tool: tool.into(),
paths: vec![],
env: vec![],
trust_level: None,
args_match: None,
capabilities: vec![],
})
.collect(),
..Default::default()
};
let enforcer = zeph_tools::PolicyEnforcer::compile(&policy_config).unwrap();
let policy_context = Arc::new(RwLock::new(zeph_tools::PolicyContext {
trust_level: zeph_common::SkillTrustLevel::Trusted,
env: std::collections::HashMap::new(),
}));
let gated = zeph_tools::PolicyGateExecutor::new(
zeph_tools::DynExecutor(session_composite),
Arc::new(enforcer),
policy_context,
);
for tool_id in ["load_skill", "memory_search", "write_file", "bash"] {
let call = zeph_tools::ToolCall {
tool_id: tool_id.into(),
params: serde_json::Map::new(),
caller_id: None,
context: None,
tool_call_id: String::new(),
skill_name: None,
};
let result = gated.execute_tool_call(&call).await;
assert!(
matches!(result, Err(zeph_tools::ToolError::Blocked { .. })),
"expected Blocked for {tool_id} from PolicyGateExecutor wrapping the full \
per-session composite (including ACP-native fs/shell), got {result:?}"
);
}
}
#[tokio::test]
async fn adversarial_policy_gate_denies_skill_and_memory_tools_in_full_acp_session_composite() {
struct AlwaysDenyLlm;
impl zeph_tools::PolicyLlmClient for AlwaysDenyLlm {
fn chat<'a>(
&'a self,
_messages: &'a [zeph_tools::PolicyMessage],
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<String, String>> + Send + 'a>,
> {
Box::pin(async move { Ok("DENY: test policy".to_owned()) })
}
}
let (session_composite, _trust_snapshot) =
build_full_acp_session_composite_with_native_fs_shell().await;
let validator = Arc::new(zeph_tools::PolicyValidator::new(
vec!["never allow load_skill, memory_search, write_file, or bash".to_owned()],
std::time::Duration::from_millis(500),
false,
vec![],
));
let llm_client: Arc<dyn zeph_tools::PolicyLlmClient> = Arc::new(AlwaysDenyLlm);
let gated = zeph_tools::AdversarialPolicyGateExecutor::new(
zeph_tools::DynExecutor(session_composite),
validator,
llm_client,
);
for tool_id in ["load_skill", "memory_search", "write_file", "bash"] {
let call = zeph_tools::ToolCall {
tool_id: tool_id.into(),
params: serde_json::Map::new(),
caller_id: None,
context: None,
tool_call_id: String::new(),
skill_name: None,
};
let result = gated.execute_tool_call(&call).await;
assert!(
matches!(result, Err(zeph_tools::ToolError::Blocked { .. })),
"expected Blocked for {tool_id} from AdversarialPolicyGateExecutor wrapping the \
full per-session composite (including ACP-native fs/shell), got {result:?}"
);
}
}
#[test]
fn build_acp_provider_factory_masks_when_registry_present() {
let mut config = zeph_core::config::Config::default();
config.llm.providers = vec![zeph_core::config::ProviderEntry {
provider_type: zeph_core::config::ProviderKind::Ollama,
name: Some("ollama".into()),
model: Some("qwen3:8b".into()),
..zeph_core::config::ProviderEntry::default()
}];
let registry = std::sync::Arc::new(zeph_sanitizer::secret_mask::SecretMaskRegistry::new());
let factory = build_acp_provider_factory(&config, Some(std::sync::Arc::clone(®istry)));
let provider = factory("ollama:qwen3:8b").expect("factory must resolve a known model key");
assert!(
matches!(provider, zeph_llm::any::AnyProvider::Masked(_)),
"factory output must be wrapped when a secret registry is supplied"
);
}
#[test]
fn build_acp_provider_factory_unmasked_when_registry_absent() {
let mut config = zeph_core::config::Config::default();
config.llm.providers = vec![zeph_core::config::ProviderEntry {
provider_type: zeph_core::config::ProviderKind::Ollama,
name: Some("ollama".into()),
model: Some("qwen3:8b".into()),
..zeph_core::config::ProviderEntry::default()
}];
let factory = build_acp_provider_factory(&config, None);
let provider = factory("ollama:qwen3:8b").expect("factory must resolve a known model key");
assert!(
!matches!(provider, zeph_llm::any::AnyProvider::Masked(_)),
"no registry supplied — factory output must be a plain passthrough"
);
}
#[test]
fn acp_provider_names_maps_known_protocols() {
let mut config = zeph_core::config::Config::default();
config.llm.providers = vec![
zeph_core::config::ProviderEntry {
provider_type: zeph_core::config::ProviderKind::Claude,
name: Some("claude".into()),
..zeph_core::config::ProviderEntry::default()
},
zeph_core::config::ProviderEntry {
provider_type: zeph_core::config::ProviderKind::OpenAi,
name: Some("openai".into()),
..zeph_core::config::ProviderEntry::default()
},
zeph_core::config::ProviderEntry {
provider_type: zeph_core::config::ProviderKind::Compatible,
name: Some("compat".into()),
..zeph_core::config::ProviderEntry::default()
},
zeph_core::config::ProviderEntry {
provider_type: zeph_core::config::ProviderKind::Ollama,
name: Some("ollama".into()),
..zeph_core::config::ProviderEntry::default()
},
];
let names = acp_provider_names(&config);
assert_eq!(
names,
vec![
("claude".to_owned(), zeph_acp::LlmProtocol::Anthropic),
("openai".to_owned(), zeph_acp::LlmProtocol::OpenAi),
("compat".to_owned(), zeph_acp::LlmProtocol::OpenAi),
(
"ollama".to_owned(),
zeph_acp::LlmProtocol::Other("ollama".to_owned())
),
]
);
}
#[test]
fn acp_provider_names_empty_providers_returns_empty_vec() {
let mut config = zeph_core::config::Config::default();
config.llm.providers.clear();
assert!(acp_provider_names(&config).is_empty());
}
#[test]
#[serial]
fn collect_project_rules_mixed_sources() {
let tmp = TempDir::new().unwrap();
make_rules_dir(tmp.path(), &["branching.md"]);
let skill_file = tmp.path().join("SKILL.md");
fs::write(&skill_file, b"").unwrap();
let orig = std::env::current_dir().unwrap();
std::env::set_current_dir(tmp.path()).unwrap();
let result = collect_project_rules(std::slice::from_ref(&skill_file));
std::env::set_current_dir(orig).unwrap();
assert_eq!(result.len(), 2);
let names: Vec<_> = result
.iter()
.filter_map(|p| p.file_name())
.map(|n| n.to_string_lossy().into_owned())
.collect();
assert!(names.contains(&"branching.md".to_owned()));
assert!(names.contains(&"SKILL.md".to_owned()));
}
#[test]
fn shared_agent_deps_has_document_and_graph_config_fields() {
let doc_cfg = zeph_core::config::DocumentConfig {
rag_enabled: true,
top_k: 7,
collection: String::new(),
chunk_size: 0,
chunk_overlap: 0,
};
assert!(doc_cfg.rag_enabled);
assert_eq!(doc_cfg.top_k, 7);
}
#[test]
fn shared_agent_deps_has_anomaly_and_orchestration_config_fields() {
let anomaly_cfg = zeph_tools::AnomalyConfig {
enabled: true,
..Default::default()
};
let orch_cfg = zeph_core::config::OrchestrationConfig {
enabled: true,
..Default::default()
};
assert!(anomaly_cfg.enabled);
assert!(orch_cfg.enabled);
}
#[cfg(all(feature = "acp-http", feature = "session"))]
#[tokio::test]
#[allow(clippy::too_many_lines)] async fn build_combined_deps_wires_skill_matching_config_from_config() {
let mut config =
zeph_core::config::Config::load(std::path::Path::new("/nonexistent")).unwrap();
config.llm.providers = vec![zeph_core::config::ProviderEntry {
provider_type: zeph_core::config::ProviderKind::Ollama,
base_url: Some("http://127.0.0.1:1".to_owned()),
model: Some("test-model".to_owned()),
..Default::default()
}];
config.memory.sqlite_path = ":memory:".to_owned();
config.skills.disambiguation_threshold = 0.55;
config.skills.two_stage_matching = true;
config.skills.confusability_threshold = 0.65;
config.skills.group_structured = true;
config.skills.support_similarity_threshold = 0.73;
config.skills.min_injection_score = 0.35;
config.skills.generation_provider = zeph_common::ProviderName::new("gen-test");
config.skills.disambiguate_provider = zeph_common::ProviderName::new("disamb-test");
config.skills.semantic_scan = true;
config.skills.semantic_scan_provider = zeph_common::ProviderName::new("scan-test");
config.skills.trust.default_level = zeph_common::SkillTrustLevel::Quarantined;
config.skills.trust.local_level = zeph_common::SkillTrustLevel::Trusted;
config.skills.rl_routing_enabled = true;
config.skills.rl_learning_rate = 0.05;
config.skills.rl_weight = 0.3;
config.skills.rl_persist_interval = 5;
config.skills.rl_warmup_updates = 3;
config.skills.rl_embed_dim = Some(8);
let app = crate::bootstrap::AppBuilder::for_test(config);
let cancel = tokio_util::sync::CancellationToken::new();
let supervisor = std::sync::Arc::new(zeph_common::TaskSupervisor::new(cancel));
let (serve_deps, acp_deps, _keepalive) = Box::pin(build_combined_deps(&app, &supervisor))
.await
.expect("build_combined_deps must succeed against a mock-provider AppBuilder");
assert!(
(serve_deps.skill_disambiguation_threshold - 0.55).abs() < f32::EPSILON,
"config.skills.disambiguation_threshold must flow into ServeAgentDeps"
);
assert!(
serve_deps.skill_two_stage_matching,
"config.skills.two_stage_matching must flow into ServeAgentDeps"
);
assert!(
(serve_deps.skill_confusability_threshold - 0.65).abs() < f32::EPSILON,
"config.skills.confusability_threshold must flow into ServeAgentDeps"
);
assert!(
serve_deps.skill_group_structured,
"config.skills.group_structured must flow into ServeAgentDeps"
);
assert!(
(serve_deps.skill_support_similarity_threshold - 0.73).abs() < f32::EPSILON,
"config.skills.support_similarity_threshold must flow into ServeAgentDeps"
);
assert!(
(serve_deps.skill_min_injection_score - 0.35).abs() < f32::EPSILON,
"config.skills.min_injection_score must flow into ServeAgentDeps"
);
assert_eq!(serve_deps.skill_generation_provider, "gen-test");
assert_eq!(serve_deps.skill_disambiguate_provider, "disamb-test");
assert!(
serve_deps.semantic_scan,
"config.skills.semantic_scan must flow into ServeAgentDeps"
);
assert_eq!(serve_deps.semantic_scan_provider, "scan-test");
assert_eq!(
serve_deps.trust_config.default_level,
zeph_common::SkillTrustLevel::Quarantined,
"config.skills.trust.default_level must flow into ServeAgentDeps"
);
assert_eq!(
serve_deps.trust_config.local_level,
zeph_common::SkillTrustLevel::Trusted,
"config.skills.trust.local_level must flow into ServeAgentDeps"
);
assert!(
serve_deps.rl_routing_enabled,
"config.skills.rl_routing_enabled must flow into ServeAgentDeps"
);
assert!(
(serve_deps.rl_learning_rate - 0.05).abs() < f32::EPSILON,
"config.skills.rl_learning_rate must flow into ServeAgentDeps"
);
assert!(
(serve_deps.rl_weight - 0.3).abs() < f32::EPSILON,
"config.skills.rl_weight must flow into ServeAgentDeps"
);
assert_eq!(
serve_deps.rl_persist_interval, 5,
"config.skills.rl_persist_interval must flow into ServeAgentDeps"
);
assert_eq!(
serve_deps.rl_warmup_updates, 3,
"config.skills.rl_warmup_updates must flow into ServeAgentDeps"
);
let serve_rl_head = serve_deps
.rl_head
.clone()
.expect("rl_head must be Some when rl_routing_enabled and rl_embed_dim resolves");
assert_eq!(
serve_rl_head.embed_dim(),
8,
"the resolved RL embed dim (config.skills.rl_embed_dim) must flow into the \
SharedCore::rl_head loaded for ServeAgentDeps"
);
assert!(
(acp_deps.skill_disambiguation_threshold - 0.55).abs() < f32::EPSILON,
"config.skills.disambiguation_threshold must flow into SharedAgentDeps"
);
assert!(
acp_deps.skill_two_stage_matching,
"config.skills.two_stage_matching must flow into SharedAgentDeps"
);
assert!(
(acp_deps.skill_confusability_threshold - 0.65).abs() < f32::EPSILON,
"config.skills.confusability_threshold must flow into SharedAgentDeps"
);
assert!(
acp_deps.skill_group_structured,
"config.skills.group_structured must flow into SharedAgentDeps"
);
assert!(
(acp_deps.skill_support_similarity_threshold - 0.73).abs() < f32::EPSILON,
"config.skills.support_similarity_threshold must flow into SharedAgentDeps"
);
assert!(
(acp_deps.skill_min_injection_score - 0.35).abs() < f32::EPSILON,
"config.skills.min_injection_score must flow into SharedAgentDeps"
);
assert_eq!(acp_deps.skill_generation_provider, "gen-test");
assert_eq!(acp_deps.skill_disambiguate_provider, "disamb-test");
assert!(
acp_deps.semantic_scan,
"config.skills.semantic_scan must flow into SharedAgentDeps"
);
assert_eq!(acp_deps.semantic_scan_provider, "scan-test");
assert_eq!(
acp_deps.trust_config.default_level,
zeph_common::SkillTrustLevel::Quarantined,
"config.skills.trust.default_level must flow into SharedAgentDeps"
);
assert_eq!(
acp_deps.trust_config.local_level,
zeph_common::SkillTrustLevel::Trusted,
"config.skills.trust.local_level must flow into SharedAgentDeps"
);
assert!(
acp_deps.rl_routing_enabled,
"config.skills.rl_routing_enabled must flow into SharedAgentDeps"
);
assert!(
(acp_deps.rl_learning_rate - 0.05).abs() < f32::EPSILON,
"config.skills.rl_learning_rate must flow into SharedAgentDeps"
);
assert!(
(acp_deps.rl_weight - 0.3).abs() < f32::EPSILON,
"config.skills.rl_weight must flow into SharedAgentDeps"
);
assert_eq!(
acp_deps.rl_persist_interval, 5,
"config.skills.rl_persist_interval must flow into SharedAgentDeps"
);
assert_eq!(
acp_deps.rl_warmup_updates, 3,
"config.skills.rl_warmup_updates must flow into SharedAgentDeps"
);
let acp_rl_head = acp_deps
.rl_head
.clone()
.expect("rl_head must be Some when rl_routing_enabled and rl_embed_dim resolves");
assert_eq!(
acp_rl_head.embed_dim(),
8,
"the resolved RL embed dim (config.skills.rl_embed_dim) must flow into the \
SharedCore::rl_head loaded for SharedAgentDeps"
);
let q = vec![0.0f32; 8];
let s = vec![0.0f32; 8];
let _ = acp_rl_head.score(&q, &s, 0.5, 0.5, 1);
assert!(acp_rl_head.update(1.0, 0.01));
assert_eq!(
serve_rl_head.update_count(),
1,
"acp_deps.rl_head and serve_deps.rl_head must share the same in-memory RoutingHead \
instance loaded once by build_shared_core (#5974)"
);
}
#[cfg(feature = "acp")]
#[tokio::test]
async fn build_acp_deps_wires_shutdown_summary_channel_identity_and_index_config_from_config() {
let mut config =
zeph_core::config::Config::load(std::path::Path::new("/nonexistent")).unwrap();
config.llm.providers = vec![zeph_core::config::ProviderEntry {
provider_type: zeph_core::config::ProviderKind::Ollama,
base_url: Some("http://127.0.0.1:1".to_owned()),
model: Some("test-model".to_owned()),
..Default::default()
}];
config.memory.sqlite_path = ":memory:".to_owned();
config.memory.shutdown_summary = true;
config.memory.shutdown_summary_min_messages = 7;
config.memory.shutdown_summary_max_messages = 42;
config.memory.shutdown_summary_timeout_secs = 9;
config.memory.shutdown_summary_provider = zeph_common::ProviderName::new("summary-test");
config.session.provider_persistence = true;
config.session.persist_provider_overrides = true;
config.index.enabled = true;
config.index.mcp_enabled = true;
let app = crate::bootstrap::AppBuilder::for_test(config);
let (deps, _keepalive) = Box::pin(build_acp_deps(&app, None, None))
.await
.expect("build_acp_deps must succeed against a mock-provider AppBuilder");
assert!(
deps.shutdown_summary,
"config.memory.shutdown_summary must flow into SharedAgentDeps"
);
assert_eq!(
deps.shutdown_summary_min_messages, 7,
"config.memory.shutdown_summary_min_messages must flow into SharedAgentDeps"
);
assert_eq!(
deps.shutdown_summary_max_messages, 42,
"config.memory.shutdown_summary_max_messages must flow into SharedAgentDeps"
);
assert_eq!(
deps.shutdown_summary_timeout_secs, 9,
"config.memory.shutdown_summary_timeout_secs must flow into SharedAgentDeps"
);
assert_eq!(
deps.shutdown_summary_provider, "summary-test",
"config.memory.shutdown_summary_provider must flow into SharedAgentDeps"
);
assert!(
deps.channel_provider_persistence,
"config.session.provider_persistence must flow into SharedAgentDeps"
);
assert!(
deps.channel_persist_provider_overrides,
"config.session.persist_provider_overrides must flow into SharedAgentDeps"
);
assert!(
deps.index_config.enabled,
"config.index.enabled must flow into SharedAgentDeps"
);
assert!(
deps.index_config.mcp_enabled,
"config.index.mcp_enabled must flow into SharedAgentDeps"
);
}
#[tokio::test]
async fn broadcast_to_mpsc_forwards_items() {
let (btx, brx) = tokio::sync::broadcast::channel::<u32>(16);
let cancel = zeph_memory::CancellationToken::new();
let mut rx = broadcast_to_mpsc(brx, cancel.clone());
btx.send(1).unwrap();
btx.send(2).unwrap();
drop(btx);
assert_eq!(rx.recv().await, Some(1));
assert_eq!(rx.recv().await, Some(2));
assert_eq!(rx.recv().await, None);
cancel.cancel();
}
#[tokio::test]
async fn broadcast_to_mpsc_cancellation_stops_task() {
let (btx, brx) = tokio::sync::broadcast::channel::<u32>(16);
let cancel = zeph_memory::CancellationToken::new();
let mut rx = broadcast_to_mpsc(brx, cancel.clone());
cancel.cancel();
tokio::task::yield_now().await;
drop(btx);
assert_eq!(rx.recv().await, None);
}
#[tokio::test]
async fn broadcast_lag_does_not_block_direct_cancel_signal() {
let (btx, brx) = tokio::sync::broadcast::channel::<u32>(1);
let adapter_cancel = zeph_memory::CancellationToken::new();
let mut rx = broadcast_to_mpsc(brx, adapter_cancel.clone());
let cancel_signal = std::sync::Arc::new(tokio::sync::Notify::new());
{
let cancel_signal = std::sync::Arc::clone(&cancel_signal);
let adapter_cancel = adapter_cancel.clone();
tokio::spawn(async move {
cancel_signal.notified().await;
adapter_cancel.cancel();
});
}
btx.send(1).unwrap();
btx.send(2).unwrap();
btx.send(3).unwrap();
tokio::task::yield_now().await;
cancel_signal.notify_one();
drop(btx);
tokio::time::timeout(
std::time::Duration::from_secs(1),
adapter_cancel.cancelled(),
)
.await
.expect("direct ACP cancel signal should not be blocked by reload lag");
loop {
let next = tokio::time::timeout(std::time::Duration::from_secs(1), rx.recv())
.await
.expect("adapter receiver should shut down promptly after cancel");
if next.is_none() {
break;
}
}
}
#[tokio::test]
async fn notify_lock_degraded_falls_back_to_channel_send_status_without_notifier() {
let (mut channel, mut handle) = zeph_core::channel::LoopbackChannel::pair(8);
notify_lock_degraded(None, &mut channel).await;
let event = handle
.output_rx
.recv()
.await
.expect("channel must receive a status event");
match event {
zeph_core::LoopbackEvent::Status(text) => {
assert_eq!(text, SESSION_LOCK_DEGRADED_MESSAGE);
}
other => panic!("expected LoopbackEvent::Status, got {other:?}"),
}
}
#[tokio::test]
async fn already_locked_session_log_notifies_client_proactively_without_prompt() {
let tmp = TempDir::new().unwrap();
let session_path = tmp.path().join("already-locked-session");
let _held_lock = zeph_session::SessionEventLog::open_exclusive(&session_path)
.await
.expect("first open_exclusive must succeed and hold the lock");
let (mut channel, _handle) = zeph_core::channel::LoopbackChannel::pair(8);
let (notify_tx, mut notify_rx) = tokio::sync::mpsc::channel(8);
let session_id =
agent_client_protocol::schema::v1::SessionId::new("already-locked-test".to_owned());
let status_notifier = Some(zeph_acp::SessionStatusNotifier::new(
notify_tx,
session_id.clone(),
));
let log = open_session_log_or_notify_locked(
&session_path,
status_notifier.as_ref(),
&mut channel,
)
.await;
assert!(
log.is_none(),
"AlreadyLocked must degrade to no persistence, not fail session creation"
);
let (notification, _ack) = notify_rx.try_recv().expect(
"client must be notified proactively — synchronously, with no prompt drain needed",
);
assert_eq!(notification.session_id, session_id);
match notification.update {
agent_client_protocol::schema::v1::SessionUpdate::AgentThoughtChunk(chunk) => {
match chunk.content {
agent_client_protocol::schema::v1::ContentBlock::Text(t) => {
assert_eq!(t.text, SESSION_LOCK_DEGRADED_MESSAGE);
}
other => panic!("expected ContentBlock::Text, got {other:?}"),
}
}
other => panic!("expected AgentThoughtChunk, got {other:?}"),
}
}
fn build_acp_agent_test_embed_fn(text: &str) -> zeph_skills::matcher::EmbedFuture {
let _ = text;
Box::pin(async { Ok(vec![1.0_f32, 0.0]) })
}
#[tokio::test]
#[allow(clippy::too_many_lines)] async fn build_acp_agent_wires_skill_matching_config() {
use zeph_commands::traits::agent::AgentAccess as _;
let mut config = zeph_core::config::Config::default();
config.skills.disambiguation_threshold = 0.77;
config.skills.two_stage_matching = true;
config.skills.confusability_threshold = 0.42;
let skill_meta = zeph_skills::loader::SkillMeta {
name: "solo-skill".to_owned(),
description: "a lone skill with no confusable sibling".to_owned(),
..Default::default()
};
let inner_matcher =
zeph_skills::matcher::SkillMatcher::new(&[&skill_meta], build_acp_agent_test_embed_fn)
.await
.expect("single-skill matcher construction must succeed with a constant embed_fn");
let (_reload_tx, reload_rx) = tokio::sync::mpsc::channel(1);
let (_config_reload_tx, config_reload_rx) = tokio::sync::mpsc::channel(1);
let (_shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(false);
let shell_policy_handle =
zeph_tools::ShellExecutor::new(&zeph_tools::ShellConfig::default()).policy_handle();
let session_config = zeph_core::AgentSessionConfig::from_config(&config, 4096);
let mcp_manager = Arc::new(crate::bootstrap::create_mcp_manager_with_vault(
&config, false, None,
));
let mock_provider =
zeph_llm::any::AnyProvider::Mock(zeph_llm::mock::MockProvider::default());
let params = BuildAcpAgentParams {
provider: mock_provider.clone(),
embedding_provider: mock_provider,
registry: Arc::new(RwLock::new(zeph_skills::registry::SkillRegistry::empty())),
matcher: Some(zeph_skills::matcher::SkillMatcherBackend::InMemory(
inner_matcher,
)),
max_active_skills: 5,
tool_executor: zeph_tools::DynExecutor(Arc::new(zeph_tools::SetCwdExecutor::new(
vec![],
))),
clock: Arc::new(zeph_common::SystemClock),
session_config,
skill_disambiguation_threshold: config.skills.disambiguation_threshold,
skill_two_stage_matching: config.skills.two_stage_matching,
skill_confusability_threshold: config.skills.confusability_threshold,
skill_group_structured: config.skills.group_structured,
skill_support_similarity_threshold: config.skills.support_similarity_threshold,
skill_min_injection_score: config.skills.min_injection_score,
skill_generation_provider: config.skills.generation_provider.as_str().to_owned(),
skill_disambiguate_provider: config.skills.disambiguate_provider.as_str().to_owned(),
semantic_scan: config.skills.semantic_scan,
semantic_scan_provider: config.skills.semantic_scan_provider.as_str().to_owned(),
trust_config: config.skills.trust.clone(),
trust_snapshot: Arc::new(RwLock::new(std::collections::HashMap::new())),
quality_pipeline: None,
rl_routing_enabled: config.skills.rl_routing_enabled,
rl_learning_rate: config.skills.rl_learning_rate,
rl_weight: config.skills.rl_weight,
rl_persist_interval: config.skills.rl_persist_interval,
rl_warmup_updates: config.skills.rl_warmup_updates,
working_dir: PathBuf::from("."),
skill_paths: Vec::new(),
reload_rx,
plugin_dirs_supplier: || Vec::<PathBuf>::new(),
shutdown_rx,
config_path: PathBuf::new(),
config_reload_rx,
startup_shell_overlay: zeph_core::ShellOverlaySnapshot {
blocked: vec![],
allowed: vec![],
},
shell_policy_handle,
mcp_tools: Vec::new(),
mcp_registry: None,
mcp_manager,
mcp_shared_tools: Arc::new(RwLock::new(Vec::new())),
mcp_config: zeph_core::config::McpConfig::default(),
focus_config: zeph_core::config::FocusConfig::default(),
sidequest_config: zeph_core::config::SidequestConfig::default(),
trajectory_config: zeph_core::config::TrajectoryConfig::default(),
category_config: zeph_core::config::CategoryConfig::default(),
provider_pool: Vec::new(),
provider_config_snapshot: zeph_core::ProviderConfigSnapshot::default(),
shutdown_summary: false,
shutdown_summary_min_messages: 0,
shutdown_summary_max_messages: 0,
shutdown_summary_timeout_secs: 0,
shutdown_summary_provider: String::new(),
channel_provider_persistence: false,
channel_persist_provider_overrides: false,
safe_mode: false,
cwd_allowed_paths: Vec::new(),
tools_enabled: true,
tool_filter_config: zeph_core::config::ToolFilterConfig::default(),
};
let (channel, _handle) = zeph_core::LoopbackChannel::pair(8);
let mut agent = Box::pin(build_acp_agent(params, channel)).await;
let output = agent
.handle_skills("confusability")
.await
.expect("handle_skills(\"confusability\") must not error");
assert!(
output.contains("above 0.42"),
"config.skills.confusability_threshold = 0.42 must reach the built Agent's \
ConfusabilityReport exactly (not e.g. 0.77, disambiguation_threshold's value, from a \
swapped SkillConfigParams field); got: {output}"
);
}
}