use std::path::Path;
use std::sync::Arc;
use std::{cmp, collections::HashSet};
use crate::{AgentFactory, AgentToolDispatcher, Config, HookEngine, HooksConfig};
#[cfg(feature = "comms")]
use crate::{CommsRuntime, CoreCommsConfig};
#[cfg(feature = "comms")]
use meerkat_core::CommsRuntimeMode;
use meerkat_core::ops_lifecycle::OpsLifecycleRegistry;
use meerkat_core::{AgentEvent, format_verbose_event};
use meerkat_hooks::DefaultHookEngine;
use meerkat_tools::builtin::shell::ShellConfig;
use meerkat_tools::{
BuiltinToolConfig, CompositeDispatcherError, FileTaskStore, MemoryTaskStore, ensure_rkat_dir,
find_project_root,
};
use tokio::sync::mpsc;
use tokio::task::JoinHandle;
#[cfg(all(feature = "comms", not(target_arch = "wasm32")))]
fn canonical_session_comms_identity_root(
user_config_root: Option<&Path>,
) -> Result<std::path::PathBuf, String> {
if let Some(root) = user_config_root {
return Ok(root.join(".rkat").join("session_comms_identity"));
}
#[cfg(target_os = "macos")]
{
let home = std::env::var_os("HOME").ok_or_else(|| {
"HOME is not set; cannot resolve session comms identity root".to_string()
})?;
#[allow(clippy::needless_return)]
return Ok(std::path::PathBuf::from(home)
.join("Library")
.join("Application Support")
.join("meerkat")
.join("session_comms_identity"));
}
#[cfg(windows)]
{
let local_app_data = std::env::var_os("LOCALAPPDATA").ok_or_else(|| {
"LOCALAPPDATA is not set; cannot resolve session comms identity root".to_string()
})?;
return Ok(std::path::PathBuf::from(local_app_data)
.join("meerkat")
.join("session_comms_identity"));
}
#[cfg(all(unix, not(target_os = "macos")))]
{
if let Some(xdg_state_home) = std::env::var_os("XDG_STATE_HOME") {
return Ok(std::path::PathBuf::from(xdg_state_home)
.join("meerkat")
.join("session_comms_identity"));
}
let home = std::env::var_os("HOME").ok_or_else(|| {
"HOME is not set; cannot resolve session comms identity root".to_string()
})?;
Ok(std::path::PathBuf::from(home)
.join(".local")
.join("state")
.join("meerkat")
.join("session_comms_identity"))
}
#[cfg(not(any(target_os = "macos", windows, unix)))]
{
Err("session-scoped comms identity root is unsupported on this platform".to_string())
}
}
pub async fn resolve_layered_hooks_config(
context_root: Option<&Path>,
user_config_root: Option<&Path>,
active_config: &Config,
) -> HooksConfig {
let mut user_entries = Vec::new();
let mut context_entries = Vec::new();
if let Some(user_root) = user_config_root {
let user_cfg_path = user_root.join(".rkat").join("config.toml");
if let Ok(Some(cfg)) = read_hooks_config_from(&user_cfg_path).await {
user_entries = cfg.entries;
}
}
if let Some(context) = context_root {
let project_cfg_path = context.join(".rkat").join("config.toml");
if let Ok(Some(cfg)) = read_hooks_config_from(&project_cfg_path).await {
context_entries = cfg.entries;
}
}
let active_hooks = &active_config.hooks;
let mut layered = HooksConfig {
default_timeout_ms: active_hooks.default_timeout_ms,
payload_max_bytes: active_hooks.payload_max_bytes,
background_max_concurrency: cmp::max(1, active_hooks.background_max_concurrency),
..HooksConfig::default()
};
let mut seen_ids: HashSet<_> = HashSet::new();
for entry in &active_hooks.entries {
if seen_ids.insert(entry.id.clone()) {
layered.entries.push(entry.clone());
}
}
for entry in &context_entries {
if seen_ids.insert(entry.id.clone()) {
layered.entries.push(entry.clone());
}
}
for entry in &user_entries {
if seen_ids.insert(entry.id.clone()) {
layered.entries.push(entry.clone());
}
}
layered
}
async fn read_hooks_config_from(path: &Path) -> Result<Option<HooksConfig>, std::io::Error> {
if !tokio::fs::try_exists(path).await? {
return Ok(None);
}
let mut parsed = Config::default();
let path_buf = path.to_path_buf();
match parsed.merge_file(&path_buf).await {
Ok(()) => Ok(Some(parsed.hooks)),
Err(err) => {
tracing::warn!(
"Failed to parse hooks config at {}: {}",
path.display(),
err
);
Ok(None)
}
}
}
pub fn create_default_hook_engine(hooks_config: HooksConfig) -> Option<Arc<dyn HookEngine>> {
if hooks_config.entries.is_empty() {
return None;
}
Some(Arc::new(DefaultHookEngine::new(hooks_config)))
}
pub async fn create_dispatcher_with_builtins(
factory: &AgentFactory,
config: BuiltinToolConfig,
shell_config: Option<ShellConfig>,
external: Option<Arc<dyn AgentToolDispatcher>>,
session_id: Option<String>,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
create_dispatcher_with_builtins_with_ops_lifecycle(
factory,
config,
shell_config,
external,
session_id,
None,
)
.await
}
pub async fn create_dispatcher_with_builtins_with_ops_lifecycle(
factory: &AgentFactory,
config: BuiltinToolConfig,
shell_config: Option<ShellConfig>,
external: Option<Arc<dyn AgentToolDispatcher>>,
session_id: Option<String>,
ops_lifecycle: Option<Arc<dyn OpsLifecycleRegistry>>,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
let store = Arc::new(MemoryTaskStore::new());
factory
.build_builtin_dispatcher(
store,
config,
factory.project_root.clone(),
shell_config,
external,
session_id,
ops_lifecycle,
)
.await
}
pub async fn create_dispatcher_with_builtins_persisted(
factory: &AgentFactory,
config: BuiltinToolConfig,
shell_config: Option<ShellConfig>,
external: Option<Arc<dyn AgentToolDispatcher>>,
session_id: Option<String>,
task_store_path: impl AsRef<Path>,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
create_dispatcher_with_builtins_persisted_with_ops_lifecycle(
factory,
config,
shell_config,
external,
session_id,
task_store_path,
None,
)
.await
}
pub async fn create_dispatcher_with_builtins_persisted_with_ops_lifecycle(
factory: &AgentFactory,
config: BuiltinToolConfig,
shell_config: Option<ShellConfig>,
external: Option<Arc<dyn AgentToolDispatcher>>,
session_id: Option<String>,
task_store_path: impl AsRef<Path>,
ops_lifecycle: Option<Arc<dyn OpsLifecycleRegistry>>,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
let store = Arc::new(FileTaskStore::new(task_store_path.as_ref().to_path_buf()));
factory
.build_builtin_dispatcher(
store,
config,
factory.project_root.clone(),
shell_config,
external,
session_id,
ops_lifecycle,
)
.await
}
pub async fn create_dispatcher_with_builtins_in_project(
factory: &AgentFactory,
config: BuiltinToolConfig,
shell_config: Option<ShellConfig>,
external: Option<Arc<dyn AgentToolDispatcher>>,
session_id: Option<String>,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
create_dispatcher_with_builtins_in_project_with_ops_lifecycle(
factory,
config,
shell_config,
external,
session_id,
None,
)
.await
}
pub async fn create_dispatcher_with_builtins_in_project_with_ops_lifecycle(
factory: &AgentFactory,
config: BuiltinToolConfig,
shell_config: Option<ShellConfig>,
external: Option<Arc<dyn AgentToolDispatcher>>,
session_id: Option<String>,
ops_lifecycle: Option<Arc<dyn OpsLifecycleRegistry>>,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
let project_root_override = factory.project_root.clone();
let project_root = tokio::task::spawn_blocking(move || {
if let Some(root) = project_root_override {
ensure_rkat_dir(&root).map_err(CompositeDispatcherError::Io)?;
return Ok::<_, CompositeDispatcherError>(root);
}
let cwd = std::env::current_dir().map_err(CompositeDispatcherError::Io)?;
let project_root =
find_project_root(&cwd).ok_or_else(|| CompositeDispatcherError::ToolInitFailed {
name: "project_root".to_string(),
message: "No .rkat directory found in current or parent directories".to_string(),
})?;
ensure_rkat_dir(&project_root).map_err(CompositeDispatcherError::Io)?;
Ok(project_root)
})
.await
.map_err(|e| CompositeDispatcherError::ToolInitFailed {
name: "project_root".to_string(),
message: format!("Failed to resolve project root: {e}"),
})??;
let store = Arc::new(FileTaskStore::in_project(&project_root));
factory
.build_builtin_dispatcher(
store,
config,
Some(project_root),
shell_config,
external,
session_id,
ops_lifecycle,
)
.await
}
pub async fn create_builtins_dispatcher(
factory: &AgentFactory,
config: BuiltinToolConfig,
session_id: Option<String>,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
create_dispatcher_with_builtins(factory, config, None, None, session_id).await
}
pub async fn create_builtins_dispatcher_with_ops_lifecycle(
factory: &AgentFactory,
config: BuiltinToolConfig,
session_id: Option<String>,
ops_lifecycle: Option<Arc<dyn OpsLifecycleRegistry>>,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
create_dispatcher_with_builtins_with_ops_lifecycle(
factory,
config,
None,
None,
session_id,
ops_lifecycle,
)
.await
}
pub async fn create_shell_dispatcher(
factory: &AgentFactory,
config: BuiltinToolConfig,
shell_config: ShellConfig,
session_id: Option<String>,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
create_dispatcher_with_builtins(factory, config, Some(shell_config), None, session_id).await
}
pub async fn create_shell_dispatcher_with_ops_lifecycle(
factory: &AgentFactory,
config: BuiltinToolConfig,
shell_config: ShellConfig,
session_id: Option<String>,
ops_lifecycle: Option<Arc<dyn OpsLifecycleRegistry>>,
) -> Result<Arc<dyn AgentToolDispatcher>, CompositeDispatcherError> {
create_dispatcher_with_builtins_with_ops_lifecycle(
factory,
config,
Some(shell_config),
None,
session_id,
ops_lifecycle,
)
.await
}
#[cfg(feature = "comms")]
pub async fn build_comms_runtime_from_config(
config: &Config,
base_dir: impl AsRef<Path>,
comms_name: &str,
peer_meta: Option<meerkat_core::PeerMeta>,
) -> Result<CommsRuntime, String> {
build_comms_runtime_from_config_scoped(config, base_dir, comms_name, peer_meta, None).await
}
#[cfg(feature = "comms")]
pub async fn build_comms_runtime_from_config_scoped(
config: &Config,
base_dir: impl AsRef<Path>,
comms_name: &str,
peer_meta: Option<meerkat_core::PeerMeta>,
inproc_namespace: Option<String>,
) -> Result<CommsRuntime, String> {
build_comms_runtime_from_config_scoped_with_silent_intents(
config,
base_dir,
comms_name,
peer_meta,
inproc_namespace,
std::sync::Arc::new(std::collections::HashSet::new()),
)
.await
}
#[cfg(feature = "comms")]
#[allow(clippy::implicit_hasher)]
pub async fn build_comms_runtime_from_config_scoped_with_silent_intents(
config: &Config,
base_dir: impl AsRef<Path>,
comms_name: &str,
peer_meta: Option<meerkat_core::PeerMeta>,
inproc_namespace: Option<String>,
silent_intents: std::sync::Arc<std::collections::HashSet<String>>,
) -> Result<CommsRuntime, String> {
let event_listen_tcp = config
.comms
.event_address
.as_ref()
.map(|addr| {
addr.parse()
.map_err(|e| format!("Invalid event_address '{addr}': {e}"))
})
.transpose()?;
let runtime =
match config.comms.mode {
CommsRuntimeMode::Inproc => CommsRuntime::inproc_only_with_silent_intents(
comms_name,
inproc_namespace.clone(),
silent_intents.clone(),
)
.map_err(|e| format!("Failed to create inproc comms runtime: {e}"))?,
CommsRuntimeMode::Tcp => {
let address =
config.comms.address.as_ref().ok_or_else(|| {
"comms.address is required when comms.mode = tcp".to_string()
})?;
let listen_tcp = address
.parse()
.map_err(|e| format!("Invalid comms TCP address '{address}': {e}"))?;
let comms = CoreCommsConfig {
enabled: true,
name: comms_name.to_string(),
inproc_namespace: inproc_namespace.clone(),
listen_tcp: Some(listen_tcp),
auth: config.comms.auth,
event_listen_tcp,
..Default::default()
};
let resolved = comms.resolve_paths(base_dir.as_ref());
let mut rt = CommsRuntime::new_machine_authority_required_with_silent_intents(
resolved,
silent_intents.clone(),
)
.await
.map_err(|e| format!("Failed to create comms runtime: {e}"))?;
rt.start_listeners()
.await
.map_err(|e| format!("Failed to start comms listeners: {e}"))?;
rt
}
CommsRuntimeMode::Uds => {
let address =
config.comms.address.as_ref().ok_or_else(|| {
"comms.address is required when comms.mode = uds".to_string()
})?;
let comms = CoreCommsConfig {
enabled: true,
name: comms_name.to_string(),
inproc_namespace: inproc_namespace.clone(),
listen_uds: Some(std::path::PathBuf::from(address)),
auth: config.comms.auth,
event_listen_tcp,
..Default::default()
};
let resolved = comms.resolve_paths(base_dir.as_ref());
let mut rt = CommsRuntime::new_machine_authority_required_with_silent_intents(
resolved,
silent_intents.clone(),
)
.await
.map_err(|e| format!("Failed to create comms runtime: {e}"))?;
rt.start_listeners()
.await
.map_err(|e| format!("Failed to start comms listeners: {e}"))?;
rt
}
};
runtime.require_peer_comms_machine_authority();
if let Some(meta) = peer_meta {
runtime.set_peer_meta(meta);
}
Ok(runtime)
}
#[cfg(feature = "comms")]
#[allow(clippy::implicit_hasher, clippy::too_many_arguments)]
pub async fn build_session_scoped_comms_runtime_from_config_scoped_with_silent_intents(
config: &Config,
base_dir: impl AsRef<Path>,
user_config_root: Option<&Path>,
comms_name: &str,
peer_meta: Option<meerkat_core::PeerMeta>,
inproc_namespace: Option<String>,
session_id: &meerkat_core::SessionId,
silent_intents: std::sync::Arc<std::collections::HashSet<String>>,
session_claim_handle: std::sync::Arc<dyn meerkat_core::handles::SessionClaimHandle>,
) -> Result<CommsRuntime, String> {
let event_listen_tcp = config
.comms
.event_address
.as_ref()
.map(|addr| {
addr.parse()
.map_err(|e| format!("Invalid event_address '{addr}': {e}"))
})
.transpose()?;
let runtime = match config.comms.mode {
CommsRuntimeMode::Inproc => CommsRuntime::inproc_only_session_scoped_with_silent_intents(
comms_name,
inproc_namespace.clone(),
canonical_session_comms_identity_root(user_config_root)?,
session_id,
silent_intents.clone(),
session_claim_handle,
)
.await
.map_err(|e| format!("Failed to create inproc comms runtime: {e}"))?,
CommsRuntimeMode::Tcp => {
let address = config
.comms
.address
.as_ref()
.ok_or_else(|| "comms.address is required when comms.mode = tcp".to_string())?;
let listen_tcp = address
.parse()
.map_err(|e| format!("Invalid comms TCP address '{address}': {e}"))?;
let comms = CoreCommsConfig {
enabled: true,
name: comms_name.to_string(),
inproc_namespace: inproc_namespace.clone(),
listen_tcp: Some(listen_tcp),
auth: config.comms.auth,
event_listen_tcp,
..Default::default()
};
let resolved = comms.resolve_paths(base_dir.as_ref());
let mut rt = CommsRuntime::new_machine_authority_required_with_silent_intents(
resolved,
silent_intents.clone(),
)
.await
.map_err(|e| format!("Failed to create comms runtime: {e}"))?;
rt.start_listeners()
.await
.map_err(|e| format!("Failed to start comms listeners: {e}"))?;
rt
}
CommsRuntimeMode::Uds => {
let address = config
.comms
.address
.as_ref()
.ok_or_else(|| "comms.address is required when comms.mode = uds".to_string())?;
let comms = CoreCommsConfig {
enabled: true,
name: comms_name.to_string(),
inproc_namespace: inproc_namespace.clone(),
listen_uds: Some(std::path::PathBuf::from(address)),
auth: config.comms.auth,
event_listen_tcp,
..Default::default()
};
let resolved = comms.resolve_paths(base_dir.as_ref());
let mut rt = CommsRuntime::new_machine_authority_required_with_silent_intents(
resolved,
silent_intents.clone(),
)
.await
.map_err(|e| format!("Failed to create comms runtime: {e}"))?;
rt.start_listeners()
.await
.map_err(|e| format!("Failed to start comms listeners: {e}"))?;
rt
}
};
runtime.require_peer_comms_machine_authority();
if let Some(meta) = peer_meta {
runtime.set_peer_meta(meta);
}
Ok(runtime)
}
#[derive(Debug, Clone, Copy, Default)]
pub struct EventLoggerConfig {
pub verbose: bool,
pub stream: bool,
}
pub fn spawn_event_logger(
mut agent_event_rx: mpsc::Receiver<AgentEvent>,
config: EventLoggerConfig,
) -> JoinHandle<()> {
tokio::spawn(async move {
use std::io::Write;
while let Some(event) = agent_event_rx.recv().await {
if config.stream
&& let AgentEvent::TextDelta { delta } = &event
{
print!("{delta}");
let _ = std::io::stdout().flush();
}
if !config.verbose {
continue;
}
if let Some(line) = format_verbose_event(&event) {
eprintln!("{line}");
}
}
})
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use meerkat_core::ToolCallView;
use meerkat_core::ToolError;
use meerkat_core::ops_lifecycle::OpsLifecycleRegistry;
use meerkat_core::{
HookCapability, HookEntryConfig, HookExecutionMode, HookId, HookPoint, HookRuntimeConfig,
HookRuntimeKind,
};
use std::path::Path;
#[cfg(all(feature = "comms", not(target_arch = "wasm32")))]
fn trusted_descriptor(
name: &str,
pubkey: meerkat_comms::identity::PubKey,
address: meerkat_core::comms::PeerAddress,
) -> meerkat_core::comms::TrustedPeerDescriptor {
meerkat_core::comms::TrustedPeerDescriptor {
peer_id: pubkey.to_peer_id(),
name: meerkat_core::comms::PeerName::new(name.to_string()).expect("valid peer name"),
address,
pubkey: *pubkey.as_bytes(),
}
}
#[cfg(all(feature = "comms", not(target_arch = "wasm32")))]
fn peer_route(
name: &str,
pubkey: meerkat_comms::identity::PubKey,
) -> meerkat_core::comms::PeerRoute {
meerkat_core::comms::PeerRoute::with_display_name(
pubkey.to_peer_id(),
meerkat_core::comms::PeerName::new(name.to_string()).expect("valid peer name"),
)
}
#[cfg(all(feature = "comms", not(target_arch = "wasm32")))]
#[tokio::test]
async fn sdk_tcp_runtime_fails_closed_before_session_machine_handle() {
use meerkat_core::agent::CommsRuntime as CoreCommsRuntime;
use meerkat_core::comms::CommsCommand;
let temp = tempfile::tempdir().expect("tempdir");
let suffix = meerkat_core::SessionId::new().to_string();
let sender_name = format!("sdk-pre-authority-sender-{suffix}");
let receiver_name = format!("sdk-pre-authority-receiver-{suffix}");
let sender = CommsRuntime::inproc_only(&sender_name).expect("sender runtime");
let mut config = Config::default();
config.comms.mode = CommsRuntimeMode::Tcp;
config.comms.address = Some("127.0.0.1:0".to_string());
let receiver = build_comms_runtime_from_config(&config, temp.path(), &receiver_name, None)
.await
.expect("sdk tcp runtime");
assert!(
receiver.peer_comms_machine_authority_required(),
"public SDK listener runtimes must fail closed before a session-owned build installs machine authority"
);
assert!(
receiver.peer_comms_handle().is_none(),
"fixture must prove the pre-handle ingress path is closed"
);
CoreCommsRuntime::add_trusted_peer(
&sender,
trusted_descriptor(
&receiver_name,
receiver.public_key(),
meerkat_core::comms::PeerAddress::new(
meerkat_core::comms::PeerTransport::Inproc,
receiver_name.clone(),
),
),
)
.await
.expect("sender trusts receiver");
CoreCommsRuntime::add_trusted_peer(
&receiver,
trusted_descriptor(
&sender_name,
sender.public_key(),
meerkat_core::comms::PeerAddress::new(
meerkat_core::comms::PeerTransport::Inproc,
sender_name.clone(),
),
),
)
.await
.expect("receiver trusts sender");
let result = CoreCommsRuntime::send(
&sender,
CommsCommand::PeerMessage {
blocks: None,
to: peer_route(&receiver_name, receiver.public_key()),
body: "must not pass sdk compatibility classifier".to_string(),
handling_mode: meerkat_core::types::HandlingMode::Queue,
},
)
.await;
assert!(matches!(
result,
Err(meerkat_core::comms::SendError::AdmissionDropped {
reason: meerkat_core::comms::AdmissionDropReason::ClassificationRejected
})
));
assert!(
CoreCommsRuntime::drain_inbox_interactions(&receiver)
.await
.is_empty(),
"pre-handle SDK-built runtime must not enqueue peer ingress"
);
}
async fn dispatch_json(
dispatcher: &dyn AgentToolDispatcher,
name: &str,
args: serde_json::Value,
) -> Result<serde_json::Value, ToolError> {
let args_raw =
serde_json::value::RawValue::from_string(args.to_string()).expect("valid args json");
let call = ToolCallView {
id: "test-1",
name,
args: &args_raw,
};
let outcome = dispatcher.dispatch(call).await?;
let text = outcome.result.text_content();
serde_json::from_str(&text).or(Ok(serde_json::Value::String(text)))
}
async fn dispatch_outcome(
dispatcher: &dyn AgentToolDispatcher,
name: &str,
args: serde_json::Value,
) -> Result<meerkat_core::ops::ToolDispatchOutcome, ToolError> {
let args_raw =
serde_json::value::RawValue::from_string(args.to_string()).expect("valid args json");
let call = ToolCallView {
id: "test-outcome",
name,
args: &args_raw,
};
dispatcher.dispatch(call).await
}
#[tokio::test]
async fn test_builtin_tools_dispatch() {
let temp_dir = tempfile::tempdir().unwrap();
let factory = AgentFactory::new(temp_dir.path().join("sessions"));
let dispatcher = create_builtins_dispatcher(&factory, BuiltinToolConfig::default(), None)
.await
.unwrap();
let args = serde_json::json!({
"subject": "Integration test task",
"description": "Testing the builtin dispatcher"
});
let result = dispatch_json(dispatcher.as_ref(), "task_create", args).await;
assert!(result.is_ok());
let task = result.unwrap();
assert!(task.get("id").is_some());
assert_eq!(task.get("subject").unwrap(), "Integration test task");
let list_result =
dispatch_json(dispatcher.as_ref(), "task_list", serde_json::json!({})).await;
assert!(list_result.is_ok());
let list = list_result.unwrap();
assert!(list.is_array());
let tasks = list.as_array().unwrap();
assert_eq!(tasks.len(), 1);
}
#[tokio::test]
async fn test_create_dispatcher_in_project_dir() {
let temp_dir = tempfile::tempdir().unwrap();
let temp_path = temp_dir.path().to_path_buf();
std::fs::create_dir_all(temp_path.join(".rkat")).unwrap();
let factory = AgentFactory::new(temp_path.join(".rkat").join("sessions"))
.project_root(temp_path.clone());
let dispatcher = create_dispatcher_with_builtins_in_project(
&factory,
BuiltinToolConfig::default(),
None,
None,
Some("test-123".to_string()),
)
.await
.unwrap();
let tools = dispatcher.tools();
assert!(tools.iter().any(|t| t.name == "task_create"));
assert!(tools.iter().any(|t| t.name == "datetime"));
assert!(!tools.iter().any(|t| t.name == "wait"));
let _ = dispatch_json(
dispatcher.as_ref(),
"task_create",
serde_json::json!({"subject":"Test","description":"Persist"}),
)
.await
.unwrap();
let tasks_file = temp_path.join(".rkat").join("tasks.json");
assert!(tasks_file.exists(), "tasks.json should be created");
}
#[tokio::test]
async fn test_builtin_tools_in_project_dir() {
let temp_dir = tempfile::tempdir().unwrap();
let temp_path = temp_dir.path();
ensure_rkat_dir(temp_path).unwrap();
let factory = AgentFactory::new(temp_path.join(".rkat").join("sessions"));
let tasks_file = temp_path.join(".rkat").join("tasks.json");
let dispatcher = create_dispatcher_with_builtins_persisted(
&factory,
BuiltinToolConfig::default(),
None,
None,
Some("file-test-session".to_string()),
&tasks_file,
)
.await
.unwrap();
let create_result = dispatch_json(
dispatcher.as_ref(),
"task_create",
serde_json::json!({
"subject": "File store test",
"description": "Testing with real file storage"
}),
)
.await;
assert!(create_result.is_ok());
let task = create_result.unwrap();
let task_id = task.get("id").unwrap().as_str().unwrap();
assert_eq!(task.get("created_by_session").unwrap(), "file-test-session");
assert!(tasks_file.exists(), "tasks.json should be created");
let get_result = dispatch_json(
dispatcher.as_ref(),
"task_get",
serde_json::json!({"id": task_id}),
)
.await;
assert!(get_result.is_ok());
let retrieved = get_result.unwrap();
assert_eq!(retrieved.get("subject").unwrap(), "File store test");
}
#[tokio::test]
async fn test_shell_dispatcher_with_ops_lifecycle_emits_async_ops() {
let temp_dir = tempfile::tempdir().unwrap();
let factory = AgentFactory::new(temp_dir.path().join("sessions"));
let shell_config =
meerkat_tools::builtin::shell::ShellConfig::with_project_root(temp_dir.path().into());
let mut config = BuiltinToolConfig::default();
config.policy.enable.insert("shell".to_string());
config.policy.enable.insert("shell_job_cancel".to_string());
let registry: Arc<dyn OpsLifecycleRegistry> =
Arc::new(meerkat_runtime::RuntimeOpsLifecycleRegistry::new());
let dispatcher = create_shell_dispatcher_with_ops_lifecycle(
&factory,
config,
shell_config,
Some(meerkat_core::types::SessionId::new().to_string()),
Some(Arc::clone(®istry)),
)
.await
.unwrap();
let outcome = dispatch_outcome(
dispatcher.as_ref(),
"shell",
serde_json::json!({
"command": "sleep 60",
"background": true
}),
)
.await
.unwrap();
assert_eq!(
outcome.async_ops.len(),
1,
"sdk helper must pass ops registry through to built-in async tools"
);
let payload: serde_json::Value =
serde_json::from_str(&outcome.result.text_content()).expect("json result");
let _ = dispatch_json(
dispatcher.as_ref(),
"shell_job_cancel",
serde_json::json!({
"job_id": payload["job_id"].as_str().expect("job id"),
}),
)
.await
.unwrap();
}
#[test]
fn test_create_default_hook_engine_none_when_no_entries() {
assert!(create_default_hook_engine(HooksConfig::default()).is_none());
}
#[test]
fn test_create_default_hook_engine_some_when_entries_exist() {
let hooks = HooksConfig {
entries: vec![HookEntryConfig {
id: HookId::new("sdk-hook"),
point: HookPoint::TurnBoundary,
mode: HookExecutionMode::Foreground,
capability: HookCapability::Observe,
runtime: HookRuntimeConfig::new(
HookRuntimeKind::InProcess,
Some(serde_json::json!({"name":"sdk_hook"})),
)
.unwrap_or_default(),
..Default::default()
}],
..Default::default()
};
assert!(create_default_hook_engine(hooks).is_some());
}
fn mk_hook(id: &str, command: &str) -> HookEntryConfig {
HookEntryConfig {
id: HookId::new(id),
point: HookPoint::TurnBoundary,
mode: HookExecutionMode::Foreground,
capability: HookCapability::Observe,
runtime: HookRuntimeConfig::new(
HookRuntimeKind::Command,
Some(serde_json::json!({ "command": command })),
)
.unwrap_or_default(),
..Default::default()
}
}
async fn write_config_with_hooks(root: &Path, hooks: Vec<HookEntryConfig>) {
let mut cfg = Config::default();
cfg.hooks.entries = hooks;
let payload = toml::to_string(&cfg).expect("serialize config");
let dir = root.join(".rkat");
tokio::fs::create_dir_all(&dir)
.await
.expect("create .rkat dir");
tokio::fs::write(dir.join("config.toml"), payload)
.await
.expect("write config");
}
#[tokio::test]
async fn resolve_layered_hooks_respects_precedence() {
let temp = tempfile::tempdir().expect("tempdir");
let user_root = temp.path().join("user");
let context_root = temp.path().join("context");
tokio::fs::create_dir_all(&user_root)
.await
.expect("user root");
tokio::fs::create_dir_all(&context_root)
.await
.expect("context root");
write_config_with_hooks(
&user_root,
vec![mk_hook("dup", "echo user"), mk_hook("u", "echo u")],
)
.await;
write_config_with_hooks(
&context_root,
vec![mk_hook("dup", "echo context"), mk_hook("c", "echo c")],
)
.await;
let mut active = Config::default();
active.hooks.entries = vec![mk_hook("dup", "echo active"), mk_hook("a", "echo a")];
let resolved =
resolve_layered_hooks_config(Some(&context_root), Some(&user_root), &active).await;
let ids: Vec<String> = resolved.entries.iter().map(|h| h.id.0.clone()).collect();
assert_eq!(ids, vec!["dup", "a", "c", "u"]);
let first_runtime = resolved.entries[0]
.runtime
.config_value()
.expect("runtime config");
assert_eq!(first_runtime["command"], "echo active");
}
#[tokio::test]
async fn resolve_layered_hooks_without_roots_uses_active_only() {
let mut active = Config::default();
active.hooks.entries = vec![mk_hook("only-active", "echo active")];
let resolved = resolve_layered_hooks_config(None, None, &active).await;
assert_eq!(resolved.entries.len(), 1);
assert_eq!(resolved.entries[0].id.0, "only-active");
}
}