use anda_core::{
Agent, AgentContext, AgentInput, AgentOutput, AgentSet, BaseContext, BoxError, CacheExpiry,
CacheFeatures, CacheStoreFeatures, CancellationToken, CompletionFeatures, CompletionRequest,
FunctionDefinition, HttpFeatures, Json, KeysFeatures, ObjectMeta, Path, PutMode, PutResult,
RequestMeta, Resource, StateFeatures, StoreFeatures, ToolGroup, ToolInput, ToolOutput,
ToolProviderSet, ToolSet,
};
use candid::Principal;
use serde::{Serialize, de::DeserializeOwned};
use std::{
collections::{BTreeMap, BTreeSet},
future::Future,
sync::Arc,
time::Duration,
};
use super::{
base::BaseCtx,
engine::RemoteEngines,
runner::{CompletionRunner, CompletionStream},
};
use crate::{
extension::todo::todo_session,
model::{Model, Models},
subagent::{SubAgentSet, SubAgentSetManager},
};
pub static DYNAMIC_REMOTE_ENGINES: &str = "_engines";
pub static REMOTE_AGENT_PREFIX: &str = "RA_";
pub static REMOTE_TOOL_PREFIX: &str = "RT_";
pub static SUB_AGENT_PREFIX: &str = "SA_";
pub(crate) fn strip_prefix_ignore_ascii_case<'a>(name: &'a str, prefix: &str) -> Option<&'a str> {
let (head, tail) = name.split_at_checked(prefix.len())?;
head.eq_ignore_ascii_case(prefix).then_some(tail)
}
pub(crate) fn strip_routing_prefix(name: &str) -> Option<&str> {
[SUB_AGENT_PREFIX, REMOTE_TOOL_PREFIX, REMOTE_AGENT_PREFIX]
.into_iter()
.find_map(|prefix| strip_prefix_ignore_ascii_case(name, prefix))
}
pub(crate) fn agent_context_path(agent_name: &str) -> String {
format!("a_{}", agent_name.to_ascii_lowercase())
}
pub(crate) fn tool_context_path(tool_name: &str) -> String {
format!("t_{}", tool_name.to_ascii_lowercase())
}
#[derive(Clone)]
pub struct AgentCtx {
pub base: BaseCtx,
pub label: String,
pub(crate) root: BaseCtx,
pub(crate) models: Arc<Models>,
pub(crate) tools: Arc<ToolSet<BaseCtx>>,
pub(crate) tool_providers: Arc<ToolProviderSet<BaseCtx>>,
pub(crate) agents: Arc<AgentSet<AgentCtx>>,
pub(crate) subagents: Arc<SubAgentSetManager>,
}
impl AgentCtx {
pub(crate) fn new(
base: BaseCtx,
models: Arc<Models>,
tools: Arc<ToolSet<BaseCtx>>,
tool_providers: Arc<ToolProviderSet<BaseCtx>>,
agents: Arc<AgentSet<AgentCtx>>,
subagents: Arc<SubAgentSetManager>,
) -> Self {
base.set_state(crate::subagent::SubAgentScope::default());
Self {
base: base.clone(),
label: String::new(),
root: base,
models,
tools,
tool_providers,
agents,
subagents,
}
}
pub fn child(&self, agent_name: &str, agent_label: &str) -> Result<Self, BoxError> {
Ok(Self {
base: self
.base
.child(agent_name.to_string(), agent_context_path(agent_name))?,
label: agent_label.to_string(),
root: self.root.clone(),
models: self.models.clone(),
tools: self.tools.clone(),
tool_providers: self.tool_providers.clone(),
agents: self.agents.clone(),
subagents: self.subagents.clone(),
})
}
pub fn child_base(&self, tool_name: &str) -> Result<BaseCtx, BoxError> {
self.base
.child(self.base.agent.clone(), tool_context_path(tool_name))
}
pub(crate) fn child_with(
&self,
caller: Principal,
agent_name: &str,
agent_label: &str,
meta: RequestMeta,
) -> Result<Self, BoxError> {
let base = self.base.child_with(
caller,
agent_name.to_string(),
agent_context_path(agent_name),
meta,
)?;
let limits = self
.base
.get_state::<crate::subagent::SubAgentScope>()
.map(|scope| scope.limits().clone())
.unwrap_or_default();
base.set_state(crate::subagent::SubAgentScope::new(limits));
base.state
.write()
.remove::<crate::subagent::ExecutionIdentity>();
Ok(Self {
base,
label: agent_label.to_string(),
root: self.root.clone(),
models: self.models.clone(),
tools: self.tools.clone(),
tool_providers: self.tool_providers.clone(),
agents: self.agents.clone(),
subagents: self.subagents.clone(),
})
}
pub(crate) fn child_base_with(
&self,
caller: Principal,
agent_name: &str,
tool_name: &str,
meta: RequestMeta,
) -> Result<BaseCtx, BoxError> {
self.base.child_with(
caller,
agent_name.to_string(),
tool_context_path(tool_name),
meta,
)
}
pub fn with_caller(&self, caller: Principal) -> Self {
let base = self.base.with_caller(caller);
if caller != self.base.caller {
let limits = self
.base
.get_state::<crate::subagent::SubAgentScope>()
.map(|scope| scope.limits().clone())
.unwrap_or_default();
base.set_state(crate::subagent::SubAgentScope::new(limits));
base.state
.write()
.remove::<crate::subagent::ExecutionIdentity>();
}
Self {
base,
..self.clone()
}
}
async fn call_remote_agent(
&self,
target: Principal,
endpoint: &str,
mut args: AgentInput,
) -> Result<AgentOutput, BoxError> {
args.meta = Some(self.base.self_meta(target));
self.https_signed_rpc(endpoint, "agent_run", &(&args,))
.await
}
pub(crate) fn has_tool_lowercase(&self, lowercase_name: &str) -> bool {
self.tools.contains_lowercase(lowercase_name)
|| self.tool_providers.contains_lowercase(lowercase_name)
}
pub fn tool_groups(&self) -> Vec<ToolGroup> {
let mut groups = BTreeMap::new();
let mut visible_names = BTreeSet::new();
let static_names: BTreeMap<String, String> = self
.tools
.iter()
.map(|(name, tool)| (name.to_string(), tool.name()))
.collect();
visible_names.extend(static_names.keys().cloned());
for group in self.tools.groups() {
merge_visible_group(&mut groups, group, &static_names);
}
let agent_names: BTreeMap<String, String> = self
.agents
.iter()
.filter(|(name, _)| visible_names.insert(name.to_string()))
.map(|(name, agent)| (name.to_string(), agent.name()))
.collect();
for group in self.agents.groups() {
merge_visible_group(&mut groups, group, &agent_names);
}
for (_, provider) in self.tool_providers.iter() {
let mut provider_names = BTreeMap::new();
for definition in provider.definitions(None) {
let lowercase = definition.name.to_ascii_lowercase();
if visible_names.insert(lowercase.clone()) {
provider_names.insert(lowercase, definition.name);
}
}
for group in provider.groups() {
merge_visible_group(&mut groups, group, &provider_names);
}
}
groups.into_values().collect()
}
async fn dynamic_remote_engines(&self) -> Option<RemoteEngines> {
self.root
.cache_store_get::<RemoteEngines>(DYNAMIC_REMOTE_ENGINES)
.await
.ok()
.map(|(engines, _)| engines)
}
fn remote_tool_definitions_with(
&self,
dynamic: Option<&RemoteEngines>,
endpoint: Option<&str>,
names: Option<&[String]>,
) -> Vec<FunctionDefinition> {
merge_remote_definitions(
self.base.remote.tool_definitions(endpoint, names),
dynamic.map(|engines| engines.tool_definitions(endpoint, names)),
REMOTE_TOOL_PREFIX,
)
}
fn remote_agent_definitions_with(
&self,
dynamic: Option<&RemoteEngines>,
endpoint: Option<&str>,
names: Option<&[String]>,
) -> Vec<FunctionDefinition> {
merge_remote_definitions(
self.base.remote.agent_definitions(endpoint, names),
dynamic.map(|engines| engines.agent_definitions(endpoint, names)),
REMOTE_AGENT_PREFIX,
)
}
pub fn completion_iter(
self,
req: CompletionRequest,
resources: Vec<Resource>,
) -> CompletionRunner {
let label = req.model.as_deref().unwrap_or(&self.label);
let model = self
.models
.resolve(label)
.unwrap_or_else(Model::not_implemented);
todo_session(&self.base);
CompletionRunner::new(self, req, model, resources)
}
pub fn completion_stream(
self,
req: CompletionRequest,
resources: Vec<Resource>,
) -> CompletionStream {
CompletionStream::new(self.completion_iter(req, resources))
}
}
impl CacheStoreFeatures for AgentCtx {}
impl AgentContext for AgentCtx {
fn tool_definitions(&self, names: Option<&[String]>) -> Vec<FunctionDefinition> {
let mut definitions = self.tools.definitions(names);
let mut seen: BTreeSet<String> =
BTreeSet::from_iter(definitions.iter().map(|d| d.name.to_ascii_lowercase()));
for definition in self.tool_providers.definitions(names) {
if seen.insert(definition.name.to_ascii_lowercase()) {
definitions.push(definition);
}
}
definitions
}
async fn remote_tool_definitions(
&self,
endpoint: Option<&str>,
names: Option<&[String]>,
) -> Result<Vec<FunctionDefinition>, BoxError> {
if let Some(names) = names
&& names.is_empty()
{
return Ok(Vec::new());
}
let dynamic = self.dynamic_remote_engines().await;
Ok(self.remote_tool_definitions_with(dynamic.as_ref(), endpoint, names))
}
async fn select_tool_resources(
&self,
prefixed_name: &str,
resources: &mut Vec<Resource>,
) -> Vec<Resource> {
if let Some(name) = strip_prefix_ignore_ascii_case(prefixed_name, REMOTE_TOOL_PREFIX) {
let res = self.base.remote.select_tool_resources(name, resources);
if !res.is_empty() {
return res;
}
if let Some(engines) = self.dynamic_remote_engines().await {
return engines.select_tool_resources(name, resources);
}
}
if self.tools.contains(prefixed_name) {
return self.tools.select_resources(prefixed_name, resources);
}
self.tool_providers
.select_resources(prefixed_name, resources)
}
fn agent_definitions(&self, names: Option<&[String]>) -> Vec<FunctionDefinition> {
if let Some(names) = names
&& names.is_empty()
{
return Vec::new();
}
let mut defs = self.agents.definitions(names);
let sub_names = names.map(|names| {
names
.iter()
.map(|name| {
strip_prefix_ignore_ascii_case(name, SUB_AGENT_PREFIX)
.unwrap_or(name)
.to_ascii_lowercase()
})
.collect::<Vec<_>>()
});
defs.extend(
self.subagents
.definitions(sub_names.as_deref())
.into_iter()
.map(|d| d.name_with_prefix(SUB_AGENT_PREFIX)),
);
defs
}
async fn remote_agent_definitions(
&self,
endpoint: Option<&str>,
names: Option<&[String]>,
) -> Result<Vec<FunctionDefinition>, BoxError> {
if let Some(names) = names
&& names.is_empty()
{
return Ok(Vec::new());
}
let dynamic = self.dynamic_remote_engines().await;
Ok(self.remote_agent_definitions_with(dynamic.as_ref(), endpoint, names))
}
async fn select_agent_resources(
&self,
prefixed_name: &str,
resources: &mut Vec<Resource>,
) -> Vec<Resource> {
if let Some(name) = strip_prefix_ignore_ascii_case(prefixed_name, REMOTE_AGENT_PREFIX) {
let res = self.base.remote.select_agent_resources(name, resources);
if !res.is_empty() {
return res;
}
if let Some(engines) = self.dynamic_remote_engines().await {
return engines.select_agent_resources(name, resources);
}
}
if let Some(prefix) = strip_prefix_ignore_ascii_case(prefixed_name, SUB_AGENT_PREFIX) {
let res = self.subagents.select_resources(prefix, resources);
if !res.is_empty() {
return res;
}
}
self.agents.select_resources(prefixed_name, resources)
}
async fn definitions(&self, names: Option<&[String]>) -> Vec<FunctionDefinition> {
if let Some(names) = names
&& names.is_empty()
{
return Vec::new();
}
let mut definitions = Vec::new();
let mut seen: BTreeSet<String> = BTreeSet::new();
let mut extend_unique = |source: Vec<FunctionDefinition>, definitions: &mut Vec<_>| {
for definition in source {
if seen.insert(definition.name.to_ascii_lowercase()) {
definitions.push(definition);
}
}
};
extend_unique(self.tool_definitions(names), &mut definitions);
extend_unique(self.agent_definitions(names), &mut definitions);
let dynamic = self.dynamic_remote_engines().await;
extend_unique(
self.remote_tool_definitions_with(dynamic.as_ref(), None, names),
&mut definitions,
);
extend_unique(
self.remote_agent_definitions_with(dynamic.as_ref(), None, names),
&mut definitions,
);
definitions
}
async fn tool_call(
&self,
mut input: ToolInput<Json>,
) -> Result<(ToolOutput<Json>, Option<Principal>), BoxError> {
if let Some(name) = strip_prefix_ignore_ascii_case(&input.name, REMOTE_TOOL_PREFIX) {
if let Some((id, endpoint, tool_name)) = self.base.remote.get_tool_endpoint(name) {
input.name = tool_name;
return self
.base
.call_remote_tool(id, &endpoint, input)
.await
.map(|output| (output, Some(id)));
}
if let Some(engines) = self.dynamic_remote_engines().await
&& let Some((id, endpoint, tool_name)) = engines.get_tool_endpoint(name)
{
input.name = tool_name;
return self
.base
.call_remote_tool(id, &endpoint, input)
.await
.map(|output| (output, Some(id)));
}
}
let ctx = self.child_base(&input.name)?;
if let Some(tool) = self.tools.get(&input.name) {
return tool
.call(ctx, input.args, input.resources)
.await
.map(|output| (output, None));
}
self.tool_providers
.call(ctx, input)
.await
.map(|output| (output, None))
}
fn agent_run(
self,
mut input: AgentInput,
) -> impl Future<Output = Result<(AgentOutput, Option<Principal>), BoxError>> + Send {
let ctx = self;
Box::pin(async move {
if let Some(name) = strip_prefix_ignore_ascii_case(&input.name, REMOTE_AGENT_PREFIX) {
if let Some((id, endpoint, agent_name)) = ctx.base.remote.get_agent_endpoint(name) {
input.name = agent_name;
return ctx
.call_remote_agent(id, &endpoint, input)
.await
.map(|output| (output, Some(id)));
}
if let Some(engines) = ctx.dynamic_remote_engines().await
&& let Some((id, endpoint, agent_name)) = engines.get_agent_endpoint(name)
{
input.name = agent_name;
return ctx
.call_remote_agent(id, &endpoint, input)
.await
.map(|output| (output, Some(id)));
}
}
if let Some(name) = strip_prefix_ignore_ascii_case(&input.name, SUB_AGENT_PREFIX) {
let name = name.to_ascii_lowercase();
if let Some(agent) = ctx.subagents.get_lowercase(&name) {
let child = ctx.child(&name, &name)?;
return agent
.run(child, input.prompt, input.resources)
.await
.map(|output| (output, None));
}
}
let name = input.name.to_ascii_lowercase();
if let Some(agent) = ctx.agents.get(&name) {
let child = ctx.child(&name, agent.label())?;
agent
.run(child, input.prompt, input.resources)
.await
.map(|output| (output, None))
} else {
Err(format!("agent {} not found", name).into())
}
})
}
async fn remote_agent_run(
&self,
endpoint: &str,
args: AgentInput,
) -> Result<AgentOutput, BoxError> {
let target = self.base.remote_target(endpoint).await?;
self.call_remote_agent(target, endpoint, args).await
}
}
impl CompletionFeatures for AgentCtx {
fn model_name(&self) -> String {
self.models
.get_model()
.unwrap_or_else(Model::not_implemented)
.model_name()
}
fn completion(
&self,
req: CompletionRequest,
resources: Vec<Resource>,
) -> impl Future<Output = Result<AgentOutput, BoxError>> + Send {
let ctx = self.clone();
Box::pin(async move { ctx.completion_iter(req, resources).run_to_end().await })
}
}
impl BaseContext for AgentCtx {
async fn remote_tool_call(
&self,
endpoint: &str,
args: ToolInput<Json>,
) -> Result<ToolOutput<Json>, BoxError> {
self.base.remote_tool_call(endpoint, args).await
}
}
impl StateFeatures for AgentCtx {
fn engine_id(&self) -> &Principal {
&self.base.id
}
fn engine_name(&self) -> &str {
&self.base.name
}
fn caller(&self) -> &Principal {
&self.base.caller
}
fn meta(&self) -> &RequestMeta {
&self.base.meta
}
fn cancellation_token(&self) -> CancellationToken {
self.base.cancellation_token.clone()
}
fn time_elapsed(&self) -> Duration {
self.base.time_elapsed()
}
}
impl KeysFeatures for AgentCtx {
async fn a256gcm_key(&self, derivation_path: Vec<Vec<u8>>) -> Result<[u8; 32], BoxError> {
self.base.a256gcm_key(derivation_path).await
}
async fn ed25519_sign_message(
&self,
derivation_path: Vec<Vec<u8>>,
message: &[u8],
) -> Result<[u8; 64], BoxError> {
self.base
.ed25519_sign_message(derivation_path, message)
.await
}
async fn ed25519_verify(
&self,
derivation_path: Vec<Vec<u8>>,
message: &[u8],
signature: &[u8],
) -> Result<(), BoxError> {
self.base
.ed25519_verify(derivation_path, message, signature)
.await
}
async fn ed25519_public_key(
&self,
derivation_path: Vec<Vec<u8>>,
) -> Result<[u8; 32], BoxError> {
self.base.ed25519_public_key(derivation_path).await
}
async fn secp256k1_sign_message_bip340(
&self,
derivation_path: Vec<Vec<u8>>,
message: &[u8],
) -> Result<[u8; 64], BoxError> {
self.base
.secp256k1_sign_message_bip340(derivation_path, message)
.await
}
async fn secp256k1_verify_bip340(
&self,
derivation_path: Vec<Vec<u8>>,
message: &[u8],
signature: &[u8],
) -> Result<(), BoxError> {
self.base
.secp256k1_verify_bip340(derivation_path, message, signature)
.await
}
async fn secp256k1_sign_message_ecdsa(
&self,
derivation_path: Vec<Vec<u8>>,
message: &[u8],
) -> Result<[u8; 64], BoxError> {
self.base
.secp256k1_sign_message_ecdsa(derivation_path, message)
.await
}
async fn secp256k1_sign_digest_ecdsa(
&self,
derivation_path: Vec<Vec<u8>>,
message_hash: &[u8],
) -> Result<[u8; 64], BoxError> {
self.base
.secp256k1_sign_digest_ecdsa(derivation_path, message_hash)
.await
}
async fn secp256k1_verify_ecdsa(
&self,
derivation_path: Vec<Vec<u8>>,
message_hash: &[u8],
signature: &[u8],
) -> Result<(), BoxError> {
self.base
.secp256k1_verify_ecdsa(derivation_path, message_hash, signature)
.await
}
async fn secp256k1_public_key(
&self,
derivation_path: Vec<Vec<u8>>,
) -> Result<[u8; 33], BoxError> {
self.base.secp256k1_public_key(derivation_path).await
}
}
impl StoreFeatures for AgentCtx {
async fn store_get(&self, path: &Path) -> Result<(bytes::Bytes, ObjectMeta), BoxError> {
self.base.store_get(path).await
}
async fn store_list(
&self,
prefix: Option<&Path>,
offset: &Path,
) -> Result<Vec<ObjectMeta>, BoxError> {
self.base.store_list(prefix, offset).await
}
async fn store_put(
&self,
path: &Path,
mode: PutMode,
value: bytes::Bytes,
) -> Result<PutResult, BoxError> {
self.base.store_put(path, mode, value).await
}
async fn store_rename_if_not_exists(&self, from: &Path, to: &Path) -> Result<(), BoxError> {
self.base.store_rename_if_not_exists(from, to).await
}
async fn store_delete(&self, path: &Path) -> Result<(), BoxError> {
self.base.store_delete(path).await
}
}
impl CacheFeatures for AgentCtx {
fn cache_contains(&self, key: &str) -> bool {
self.base.cache_contains(key)
}
async fn cache_get<T>(&self, key: &str) -> Result<T, BoxError>
where
T: DeserializeOwned,
{
self.base.cache_get(key).await
}
async fn cache_get_with<T, F>(&self, key: &str, init: F) -> Result<T, BoxError>
where
T: Sized + DeserializeOwned + Serialize + Send,
F: Future<Output = Result<(T, Option<CacheExpiry>), BoxError>> + Send + 'static,
{
self.base.cache_get_with(key, init).await
}
async fn cache_set<T>(&self, key: &str, val: (T, Option<CacheExpiry>))
where
T: Sized + Serialize + Send,
{
self.base.cache_set(key, val).await
}
async fn cache_set_if_not_exists<T>(&self, key: &str, val: (T, Option<CacheExpiry>)) -> bool
where
T: Sized + Serialize + Send,
{
self.base.cache_set_if_not_exists(key, val).await
}
async fn cache_delete(&self, key: &str) -> bool {
self.base.cache_delete(key).await
}
}
impl HttpFeatures for AgentCtx {
async fn https_call(
&self,
url: &str,
method: http::Method,
headers: Option<http::HeaderMap>,
body: Option<Vec<u8>>,
) -> Result<reqwest::Response, BoxError> {
self.base.https_call(url, method, headers, body).await
}
async fn https_signed_call(
&self,
url: &str,
method: http::Method,
message_digest: [u8; 32],
headers: Option<http::HeaderMap>,
body: Option<Vec<u8>>,
) -> Result<reqwest::Response, BoxError> {
self.base
.https_signed_call(url, method, message_digest, headers, body)
.await
}
async fn https_signed_rpc<T>(
&self,
endpoint: &str,
method: &str,
args: impl Serialize + Send,
) -> Result<T, BoxError>
where
T: DeserializeOwned,
{
self.base.https_signed_rpc(endpoint, method, args).await
}
}
fn merge_remote_definitions(
mut definitions: Vec<FunctionDefinition>,
dynamic: Option<Vec<FunctionDefinition>>,
routing_prefix: &str,
) -> Vec<FunctionDefinition> {
if let Some(dynamic) = dynamic {
let mut seen: BTreeSet<String> = definitions
.iter()
.map(|definition| definition.name.to_ascii_lowercase())
.collect();
definitions.extend(
dynamic
.into_iter()
.filter(|definition| seen.insert(definition.name.to_ascii_lowercase())),
);
}
definitions
.into_iter()
.map(|definition| definition.name_with_prefix(routing_prefix))
.collect()
}
fn merge_visible_group(
groups: &mut BTreeMap<String, ToolGroup>,
mut group: ToolGroup,
visible_names: &BTreeMap<String, String>,
) {
let id = group.id.trim();
if id.is_empty() || visible_names.is_empty() {
return;
}
let mut seen_members = BTreeSet::new();
let mut members = group
.members
.into_iter()
.filter_map(|member| {
let lowercase = member.trim().to_ascii_lowercase();
if lowercase.is_empty() || !seen_members.insert(lowercase.clone()) {
return None;
}
visible_names.get(&lowercase).cloned()
})
.collect::<Vec<_>>();
if members.is_empty() {
return;
}
members.sort_by_key(|name| name.to_ascii_lowercase());
let key = id.to_ascii_lowercase();
match groups.get_mut(&key) {
Some(existing) => {
let mut existing_members = existing
.members
.iter()
.map(|member| member.to_ascii_lowercase())
.collect::<BTreeSet<_>>();
for member in members {
if existing_members.insert(member.to_ascii_lowercase()) {
existing.members.push(member);
}
}
existing
.members
.sort_by_key(|name| name.to_ascii_lowercase());
}
None => {
group.id = id.to_string();
group.members = members;
groups.insert(key, group);
}
}
}
#[cfg(test)]
mod tests {
use anda_core::{
AgentContext as _, AgentInput, BaseContext as _, BoxError, CacheFeatures as _,
CacheStoreFeatures as _, CompletionFeatures as _, HttpFeatures as _, Json,
KeysFeatures as _, Path, PutMode, StateFeatures as _, StoreFeatures as _, ToolInput,
};
use bytes::Bytes;
use candid::Principal;
use cbor2::from_slice;
use ic_cose_types::to_cbor_bytes;
use serde_json::json;
use std::sync::Arc;
use super::{
DYNAMIC_REMOTE_ENGINES, REMOTE_AGENT_PREFIX, REMOTE_TOOL_PREFIX, SUB_AGENT_PREFIX,
};
use crate::context::test_fixtures::*;
use crate::{engine::EngineBuilder, model::Model};
#[test]
fn json_in_cbor_works() {
let json = json!({
"level": "info",
"message": "Hello, world!",
"timestamp": "2021-09-01T12:00:00Z",
"data": {
"key": "value",
"number": 42,
"flag": true
}
});
let data = to_cbor_bytes(&json);
let val: serde_json::Value = from_slice(&data[..]).unwrap();
assert_eq!(json, val);
}
#[test]
fn agent_child_context_switches_agent_namespace_and_preserves_state() {
let ctx = EngineBuilder::new().mock_ctx();
ctx.base.set_state("parent-state".to_string());
let child = ctx.child("worker_agent", "Worker").unwrap();
assert_eq!(ctx.base.agent, "Mocker");
assert_eq!(child.base.agent, "worker_agent");
assert_eq!(child.label, "Worker");
assert_eq!(child.base.path.as_ref(), "a_worker_agent");
assert_eq!(
child.base.get_state::<String>().as_deref(),
Some("parent-state")
);
let tool_ctx = child.child_base("note").unwrap();
assert_eq!(tool_ctx.agent, "worker_agent");
assert_eq!(tool_ctx.path.as_ref(), "t_note");
assert_eq!(
tool_ctx.get_state::<String>().as_deref(),
Some("parent-state")
);
assert_eq!(
child.child_base("execute_kip").unwrap().path().as_ref(),
"t_execute_kip"
);
assert_eq!(
child.child_base("memory_readonly").unwrap().path().as_ref(),
"t_memory_readonly"
);
}
#[tokio::test(flavor = "current_thread")]
async fn agent_context_definitions_dynamic_resources_and_missing_runs() {
let model = Model::with_completer(Arc::new(EchoCompleter));
let ctx = EngineBuilder::new()
.with_model(model)
.register_tool(Arc::new(EchoTool))
.unwrap()
.register_agent(Arc::new(EchoAgent), None)
.unwrap()
.mock_ctx();
let empty: Vec<String> = Vec::new();
assert!(ctx.tool_definitions(Some(&empty)).is_empty());
assert!(ctx.agent_definitions(Some(&empty)).is_empty());
assert!(
ctx.remote_tool_definitions(None, Some(&empty))
.await
.unwrap()
.is_empty()
);
assert!(
ctx.remote_agent_definitions(None, Some(&empty))
.await
.unwrap()
.is_empty()
);
assert!(ctx.definitions(Some(&empty)).await.is_empty());
ctx.root
.cache_store_set(DYNAMIC_REMOTE_ENGINES, dynamic_remote_engines(), None)
.await
.unwrap();
let definitions = ctx.definitions(None).await;
assert!(definitions.iter().any(|d| d.name == "echo_tool"));
assert!(definitions.iter().any(|d| d.name == "echo_agent"));
assert!(
definitions
.iter()
.any(|d| d.name == format!("{REMOTE_TOOL_PREFIX}dyn_lookup"))
);
assert!(
definitions
.iter()
.any(|d| d.name == format!("{REMOTE_AGENT_PREFIX}dyn_chat"))
);
let mut resources = vec![resource(1, &["text"]), resource(2, &["md"])];
let selected = ctx
.select_tool_resources(&format!("{REMOTE_TOOL_PREFIX}dyn_lookup"), &mut resources)
.await;
assert_eq!(
selected
.iter()
.map(|resource| resource._id)
.collect::<Vec<_>>(),
vec![1]
);
assert_eq!(
resources
.iter()
.map(|resource| resource._id)
.collect::<Vec<_>>(),
vec![2]
);
let selected = ctx
.select_agent_resources(&format!("{REMOTE_AGENT_PREFIX}dyn_chat"), &mut resources)
.await;
assert_eq!(
selected
.iter()
.map(|resource| resource._id)
.collect::<Vec<_>>(),
vec![2]
);
assert!(resources.is_empty());
let tool_err = ctx
.tool_call(ToolInput {
name: format!("{REMOTE_TOOL_PREFIX}dyn_lookup"),
args: json!({}),
resources: Vec::new(),
meta: None,
})
.await
.unwrap_err();
assert!(tool_err.to_string().contains("not implemented"));
let agent_err = ctx
.clone()
.agent_run(AgentInput {
name: format!("{REMOTE_AGENT_PREFIX}dyn_chat"),
prompt: "hello".to_string(),
..Default::default()
})
.await
.unwrap_err();
assert!(agent_err.to_string().contains("not implemented"));
let agent_err = ctx
.clone()
.agent_run(AgentInput {
name: format!("{REMOTE_AGENT_PREFIX}missing"),
prompt: "hello".to_string(),
..Default::default()
})
.await
.unwrap_err();
assert!(agent_err.to_string().contains("agent ra_missing not found"));
let agent_err = ctx
.clone()
.agent_run(AgentInput {
name: format!("{SUB_AGENT_PREFIX}missing"),
prompt: "hello".to_string(),
..Default::default()
})
.await
.unwrap_err();
assert!(agent_err.to_string().contains("agent sa_missing not found"));
let agent_err = ctx
.agent_run(AgentInput {
name: "missing_agent".to_string(),
prompt: "hello".to_string(),
..Default::default()
})
.await
.unwrap_err();
assert!(
agent_err
.to_string()
.contains("agent missing_agent not found")
);
}
#[tokio::test(flavor = "current_thread")]
async fn agent_context_trait_forwarders_cover_base_store_cache_keys_and_http() {
let model = Model::with_completer(Arc::new(EchoCompleter));
let ctx = EngineBuilder::new().with_model(model).mock_ctx();
assert_eq!(ctx.model_name(), "echo");
assert_eq!(ctx.engine_name(), "Mocker");
assert_eq!(*ctx.engine_id(), Principal::anonymous());
assert_eq!(*ctx.caller(), Principal::anonymous());
assert!(ctx.meta().user.is_none());
assert!(!ctx.cancellation_token().is_cancelled());
assert!(ctx.time_elapsed() < std::time::Duration::from_secs(60));
let caller = Principal::self_authenticating([4; 32]);
let called_by = ctx.with_caller(caller);
assert_eq!(*called_by.caller(), caller);
let path = Path::from("agent_ctx_file");
let renamed = Path::from("agent_ctx_file_renamed");
ctx.store_put(&path, PutMode::Overwrite, Bytes::from_static(b"data"))
.await
.unwrap();
let (stored, meta) = ctx.store_get(&path).await.unwrap();
assert_eq!(stored, Bytes::from_static(b"data"));
assert_eq!(meta.location, path);
ctx.store_rename_if_not_exists(&path, &renamed)
.await
.unwrap();
let listed = ctx.store_list(None, &Path::from("")).await.unwrap();
assert!(listed.iter().any(|meta| meta.location == renamed));
ctx.store_delete(&renamed).await.unwrap();
assert!(ctx.store_get(&renamed).await.is_err());
let created: String = ctx
.cache_get_with("root_key", async {
Ok::<_, BoxError>(("created".to_string(), None))
})
.await
.unwrap();
assert_eq!(created, "created");
let cache_ctx = ctx.child("tools_search", "Tools Search").unwrap();
assert!(!cache_ctx.cache_contains("number"));
cache_ctx.cache_set("number", (42_u64, None)).await;
assert_eq!(cache_ctx.cache_get::<u64>("number").await.unwrap(), 42);
let initialized: String = cache_ctx
.cache_get_with("initialized", async {
Ok::<_, BoxError>(("created".to_string(), None))
})
.await
.unwrap();
assert_eq!(initialized, "created");
assert!(cache_ctx.cache_contains("number"));
assert!(cache_ctx.cache_delete("number").await);
assert!(!cache_ctx.cache_contains("number"));
assert!(ctx.a256gcm_key(Vec::new()).await.is_err());
assert!(ctx.ed25519_sign_message(Vec::new(), b"msg").await.is_err());
assert!(
ctx.ed25519_verify(Vec::new(), b"msg", &[0; 64])
.await
.is_err()
);
assert!(ctx.ed25519_public_key(Vec::new()).await.is_err());
assert!(
ctx.secp256k1_sign_message_bip340(Vec::new(), b"msg")
.await
.is_err()
);
assert!(
ctx.secp256k1_verify_bip340(Vec::new(), b"msg", &[0; 64])
.await
.is_err()
);
assert!(
ctx.secp256k1_sign_message_ecdsa(Vec::new(), b"msg")
.await
.is_err()
);
assert!(
ctx.secp256k1_sign_digest_ecdsa(Vec::new(), &[0; 32])
.await
.is_err()
);
assert!(
ctx.secp256k1_verify_ecdsa(Vec::new(), &[0; 32], &[0; 64])
.await
.is_err()
);
assert!(ctx.secp256k1_public_key(Vec::new()).await.is_err());
assert!(
ctx.https_call("https://example.test", http::Method::GET, None, None)
.await
.is_err()
);
assert!(
ctx.https_signed_call(
"https://example.test",
http::Method::POST,
[0; 32],
None,
Some(Vec::new()),
)
.await
.is_err()
);
let rpc: Result<Json, BoxError> = ctx
.https_signed_rpc("https://example.test", "method", &())
.await;
assert!(rpc.is_err());
let err = ctx
.remote_tool_call(
"https://missing.example",
ToolInput {
name: "lookup".to_string(),
args: json!({}),
resources: Vec::new(),
meta: None,
},
)
.await
.unwrap_err();
assert!(
err.to_string()
.contains("remote engine endpoint https://missing.example not found")
);
}
}