use async_trait::async_trait;
use chrono::Utc;
use everruns_core::ExecutionContext;
use everruns_core::MessageRetriever;
use everruns_core::capabilities::{
Capability, CapabilityStatus, SystemPromptContext, collect_capabilities_with_configs,
};
use everruns_core::session_task::{
CreateSessionTask, NewTaskMessage, SessionTask, SessionTaskFilter, SessionTaskRegistry,
SessionTaskUpdate, TaskMessage, apply_task_update, new_session_task,
};
use everruns_core::tool_context::ToolContext;
use everruns_core::{
AgentDefinition, CapabilityRegistry, EventData, ExecutionSession, HarnessDefinition,
InputMessage, SessionExecutionState, TokenUsage, Tool, ToolExecutionResult, ToolRegistry,
};
use everruns_core::{
event_emitter::EventEmitter, execution_loading::AgentStore, execution_loading::HarnessStore,
execution_loading::SessionStore, provider_resolution::ProviderStore,
session_files::SessionFileSystem,
};
use everruns_engine::{ActInput, InputAtomInput, ReasonResult};
use everruns_engine::{TurnPlan, TurnState};
use everruns_host::SessionMutator;
use everruns_host::{
InMemoryAgentStore, InMemoryHarnessStore, InMemoryProviderStore, InMemorySessionFileStore,
InProcessExecution, ResolvedTurnInputs, RuntimeHostAdapter, RuntimeSessionLifecycle,
TurnStopReason, advance_host_execution, execute_act_activity, execute_input_activity,
inspect_turn_context,
};
use everruns_provider::driver_registry::DriverRegistry;
use everruns_provider::model_spec::ModelSpec;
use everruns_provider::provider::DriverId;
use everruns_provider::tool_types::{ToolCall, ToolResult};
use everruns_provider::typed_id::{AgentId, HarnessId, MessageId, SessionId, TurnId};
use everruns_provider::user_facing_error::codes as user_facing_error_codes;
use everruns_test_support::{
InMemoryEventEmitter, InMemoryMessageRetriever, TestMathCapability,
llmsim_driver::register_driver,
};
use serde_json::json;
use std::collections::HashMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use tokio::sync::RwLock;
use uuid::Uuid;
#[derive(Clone, Default)]
struct TestSessionStore {
sessions: Arc<RwLock<HashMap<SessionId, ExecutionSession>>>,
fail_status_writes: Arc<AtomicBool>,
}
impl TestSessionStore {
async fn insert(&self, session: ExecutionSession) {
self.sessions.write().await.insert(session.id, session);
}
async fn set_status(
&self,
session_id: SessionId,
status: SessionExecutionState,
) -> everruns_provider::error::Result<ExecutionSession> {
if self.fail_status_writes.load(Ordering::SeqCst) {
return Err(everruns_provider::error::AgentLoopError::config(
"injected session status failure",
));
}
let mut sessions = self.sessions.write().await;
let session = sessions.get_mut(&session_id).expect("session exists");
session.status = status;
Ok(session.clone())
}
}
#[async_trait]
impl SessionStore for TestSessionStore {
async fn get_session(
&self,
session_id: SessionId,
) -> everruns_provider::error::Result<Option<ExecutionSession>> {
Ok(self.sessions.read().await.get(&session_id).cloned())
}
}
#[async_trait]
impl SessionMutator for TestSessionStore {
async fn update_session_title(
&self,
session_id: SessionId,
title: String,
) -> everruns_provider::error::Result<ExecutionSession> {
let mut sessions = self.sessions.write().await;
let session = sessions.get_mut(&session_id).expect("session exists");
session.title = Some(title);
Ok(session.clone())
}
}
#[derive(Clone)]
struct MockHostAdapter {
capability_registry: CapabilityRegistry,
driver_registry: DriverRegistry,
harness_store: Arc<InMemoryHarnessStore>,
agent_store: Arc<InMemoryAgentStore>,
session_store: Arc<TestSessionStore>,
message_store: Arc<InMemoryMessageRetriever>,
provider_store: Arc<InMemoryProviderStore>,
event_emitter: Arc<InMemoryEventEmitter>,
file_store: Arc<InMemorySessionFileStore>,
session_task_registry: Option<Arc<dyn SessionTaskRegistry>>,
}
#[async_trait]
impl RuntimeHostAdapter for MockHostAdapter {
async fn set_session_status(
&self,
_org_id: i64,
session_id: SessionId,
status: SessionExecutionState,
) -> everruns_provider::error::Result<()> {
self.session_store.set_status(session_id, status).await?;
Ok(())
}
async fn load_resolved_turn(
&self,
_org_id: i64,
session_id: SessionId,
) -> everruns_provider::error::Result<ResolvedTurnInputs> {
let session = self
.session_store
.get_session(session_id)
.await?
.expect("session exists");
let agent = match session.agent_id {
Some(agent_id) => self.agent_store.get_agent(agent_id).await?,
None => None,
};
let harness = self
.harness_store
.get_harness(session.harness_id)
.await?
.expect("harness exists");
let snapshot =
everruns_core::ResolvedExecutionSnapshot::project(&harness, agent.as_ref(), &session)?;
Ok(ResolvedTurnInputs {
snapshot,
messages: self.message_store.load(session_id).await?,
mcp_tool_definitions: vec![],
})
}
fn capability_registry(&self) -> CapabilityRegistry {
self.capability_registry.clone()
}
fn driver_registry(&self) -> DriverRegistry {
self.driver_registry.clone()
}
fn harness_store(&self, _org_id: i64) -> Arc<dyn HarnessStore> {
self.harness_store.clone()
}
fn agent_store(&self, _org_id: i64) -> Arc<dyn AgentStore> {
self.agent_store.clone()
}
fn session_store(&self, _org_id: i64) -> Arc<dyn SessionStore> {
self.session_store.clone()
}
fn session_mutator(&self, _org_id: i64) -> Arc<dyn SessionMutator> {
self.session_store.clone()
}
fn provider_store(&self, _org_id: i64) -> Arc<dyn ProviderStore> {
self.provider_store.clone()
}
fn message_store(&self) -> Arc<dyn everruns_core::MessageRetriever> {
self.message_store.clone()
}
fn event_emitter(&self) -> Arc<dyn EventEmitter> {
self.event_emitter.clone()
}
fn file_store(&self) -> Arc<dyn SessionFileSystem> {
self.file_store.clone()
}
fn session_task_registry(&self) -> Option<Arc<dyn SessionTaskRegistry>> {
self.session_task_registry.clone()
}
}
#[derive(Default)]
struct TestTaskRegistry {
tasks: RwLock<Vec<SessionTask>>,
}
#[async_trait]
impl SessionTaskRegistry for TestTaskRegistry {
async fn create(
&self,
input: CreateSessionTask,
) -> everruns_provider::error::Result<SessionTask> {
let task = new_session_task(input, Utc::now());
self.tasks.write().await.push(task.clone());
Ok(task)
}
async fn update(
&self,
session_id: SessionId,
task_id: &str,
update: SessionTaskUpdate,
) -> everruns_provider::error::Result<Option<SessionTask>> {
let mut tasks = self.tasks.write().await;
let Some(task) = tasks
.iter_mut()
.find(|task| task.session_id == session_id && task.id == task_id)
else {
return Ok(None);
};
apply_task_update(task, update, Utc::now());
Ok(Some(task.clone()))
}
async fn get(
&self,
session_id: SessionId,
task_id: &str,
) -> everruns_provider::error::Result<Option<SessionTask>> {
Ok(self
.tasks
.read()
.await
.iter()
.find(|task| task.session_id == session_id && task.id == task_id)
.cloned())
}
async fn list(
&self,
session_id: SessionId,
filter: Option<&SessionTaskFilter>,
) -> everruns_provider::error::Result<Vec<SessionTask>> {
Ok(self
.tasks
.read()
.await
.iter()
.filter(|task| {
task.session_id == session_id
&& filter
.and_then(|filter| filter.kind.as_deref())
.is_none_or(|kind| task.kind == kind)
})
.cloned()
.collect())
}
async fn request_cancel(
&self,
session_id: SessionId,
task_id: &str,
) -> everruns_provider::error::Result<Option<SessionTask>> {
self.get(session_id, task_id).await
}
async fn record_message(
&self,
_session_id: SessionId,
_task_id: &str,
_message: NewTaskMessage,
) -> everruns_provider::error::Result<TaskMessage> {
Err(everruns_provider::error::AgentLoopError::tool("not needed"))
}
async fn list_messages(
&self,
_session_id: SessionId,
_task_id: &str,
_limit: Option<u32>,
_after_id: Option<&str>,
) -> everruns_provider::error::Result<Vec<TaskMessage>> {
Ok(vec![])
}
}
async fn build_registry(
capability_registry: &CapabilityRegistry,
session_id: SessionId,
capabilities: &[everruns_capability::CapabilityRef],
) -> everruns_provider::error::Result<ToolRegistry> {
let ctx = SystemPromptContext::without_file_store(session_id);
let collected =
collect_capabilities_with_configs(capabilities, capability_registry, &ctx).await;
let mut registry = ToolRegistry::with_defaults();
for tool in collected.tools {
registry.register_boxed(tool);
}
Ok(registry)
}
struct OverlayEchoTool;
#[cfg(feature = "builtins")]
struct PersistedOutputTool;
#[cfg(feature = "builtins")]
#[async_trait]
impl Tool for PersistedOutputTool {
fn name(&self) -> &str {
"persisted_output_fixture"
}
fn description(&self) -> &str {
"Returns output large enough to exercise final persistence and hard limits."
}
fn parameters_schema(&self) -> serde_json::Value {
json!({ "type": "object", "additionalProperties": false })
}
fn hints(&self) -> everruns_provider::tool_types::ToolHints {
everruns_provider::tool_types::ToolHints::default().with_persist_output(true)
}
async fn execute(&self, _arguments: serde_json::Value) -> ToolExecutionResult {
let stdout = format!("{}final-marker", "x".repeat(70 * 1024));
ToolExecutionResult::success(json!({
"stdout": stdout,
"stderr": "",
"exit_code": 0,
"success": true,
}))
}
}
#[cfg(feature = "builtins")]
struct PersistedOutputCapability;
#[cfg(feature = "builtins")]
impl Capability for PersistedOutputCapability {
fn id(&self) -> &str {
"persisted_output_fixture"
}
fn name(&self) -> &str {
"Persisted Output Fixture"
}
fn description(&self) -> &str {
"Test capability for the host-owned final persistence hook."
}
fn tools(&self) -> Vec<Box<dyn Tool>> {
vec![Box::new(PersistedOutputTool)]
}
}
#[async_trait]
impl Tool for OverlayEchoTool {
fn name(&self) -> &str {
"overlay_echo"
}
fn description(&self) -> &str {
"Returns the provided value."
}
fn parameters_schema(&self) -> serde_json::Value {
json!({
"type": "object",
"properties": {
"value": {"type": "string"}
},
"required": ["value"],
"additionalProperties": false
})
}
async fn execute(&self, arguments: serde_json::Value) -> ToolExecutionResult {
ToolExecutionResult::success(json!({
"value": arguments["value"].as_str().unwrap_or_default(),
}))
}
}
struct ContextParityTool;
#[async_trait]
impl Tool for ContextParityTool {
fn name(&self) -> &str {
"context_parity"
}
fn description(&self) -> &str {
"Reports whether runtime host services reached ToolContext."
}
fn parameters_schema(&self) -> serde_json::Value {
json!({"type": "object", "additionalProperties": false})
}
async fn execute(&self, _arguments: serde_json::Value) -> ToolExecutionResult {
ToolExecutionResult::tool_error("context required")
}
async fn execute_with_context(
&self,
_arguments: serde_json::Value,
context: &ToolContext,
) -> ToolExecutionResult {
ToolExecutionResult::success(json!({
"file_store": context.file_store.is_some(),
"message_retriever": context.message_retriever.is_some(),
"session_store": context.session_store.is_some(),
"session_mutator": context.extensions.get::<everruns_host::SessionMutatorExt>().is_some(),
"agent_store": context.agent_store.is_some(),
"session_task_registry": context.session_task_registry.is_some(),
"capability_registry": context.capability_registry.is_some(),
"tool_registry": context.tool_registry.is_some(),
"org_id": context.org_id.is_some(),
"event_emitter": context.event_emitter.is_some(),
"event_context": context.event_context.is_some(),
"tool_call_id": context.tool_call_id.is_some(),
}))
}
}
struct ContextParityCapability;
impl Capability for ContextParityCapability {
fn id(&self) -> &str {
"context_parity"
}
fn name(&self) -> &str {
"Context Parity"
}
fn description(&self) -> &str {
"Test capability for runtime-owned ToolContext service propagation."
}
fn tools(&self) -> Vec<Box<dyn Tool>> {
vec![Box::new(ContextParityTool)]
}
}
struct OverlayHook;
#[async_trait]
impl everruns_core::tool_hooks::PostToolExecHook for OverlayHook {
async fn after_exec(
&self,
_tool_call: &ToolCall,
_tool_def: &everruns_provider::tool_types::ToolDefinition,
result: &mut ToolResult,
_context: &ToolContext,
) {
if let Some(value) = result
.result
.as_mut()
.and_then(|value| value.as_object_mut())
{
value.insert("hooked".to_string(), json!(true));
}
}
}
struct GateCapability;
impl Capability for GateCapability {
fn id(&self) -> &str {
"gate"
}
fn name(&self) -> &str {
"Gate"
}
fn description(&self) -> &str {
"Test capability contributing a blocking pre-tool hook."
}
fn status(&self) -> CapabilityStatus {
CapabilityStatus::Available
}
fn pre_tool_use_hooks(&self) -> Vec<Arc<dyn everruns_core::tool_hooks::PreToolUseHook>> {
vec![Arc::new(BlockEchoHook)]
}
}
struct ComingSoonGateCapability;
impl Capability for ComingSoonGateCapability {
fn id(&self) -> &str {
"coming_soon_gate"
}
fn name(&self) -> &str {
"Coming Soon Gate"
}
fn description(&self) -> &str {
"Unavailable test capability contributing a blocking pre-tool hook."
}
fn status(&self) -> CapabilityStatus {
CapabilityStatus::ComingSoon
}
fn pre_tool_use_hooks(&self) -> Vec<Arc<dyn everruns_core::tool_hooks::PreToolUseHook>> {
vec![Arc::new(BlockEchoHook)]
}
}
struct BlockEchoHook;
#[async_trait]
impl everruns_core::tool_hooks::PreToolUseHook for BlockEchoHook {
async fn before_exec(
&self,
tool_call: ToolCall,
_tool_def: &everruns_provider::tool_types::ToolDefinition,
_context: &ToolContext,
) -> everruns_core::tool_hooks::PreToolUseDecision {
if tool_call.name == "overlay_echo" {
everruns_core::tool_hooks::PreToolUseDecision::Block {
tool_call,
reason: "denied by test policy".into(),
user_message: None,
}
} else {
everruns_core::tool_hooks::PreToolUseDecision::Continue(tool_call)
}
}
}
struct OverlayEchoCapability;
impl Capability for OverlayEchoCapability {
fn id(&self) -> &str {
"overlay_echo"
}
fn name(&self) -> &str {
"Overlay Echo"
}
fn description(&self) -> &str {
"Test capability that contributes a tool and a post-tool hook."
}
fn status(&self) -> CapabilityStatus {
CapabilityStatus::Available
}
fn tools(&self) -> Vec<Box<dyn Tool>> {
vec![Box::new(OverlayEchoTool)]
}
fn post_tool_exec_hooks(&self) -> Vec<Arc<dyn everruns_core::tool_hooks::PostToolExecHook>> {
vec![Arc::new(OverlayHook)]
}
}
struct OverlayAliasCapability;
impl Capability for OverlayAliasCapability {
fn id(&self) -> &str {
"overlay_alias"
}
fn name(&self) -> &str {
"Overlay Alias"
}
fn description(&self) -> &str {
"Test capability that depends on overlay_echo."
}
fn status(&self) -> CapabilityStatus {
CapabilityStatus::Available
}
fn dependencies(&self) -> Vec<&'static str> {
vec!["overlay_echo"]
}
}
struct NarratingTool;
#[async_trait]
impl Tool for NarratingTool {
fn name(&self) -> &str {
"narrating_tool"
}
fn display_name(&self) -> Option<&str> {
Some("Narrating Tool")
}
fn description(&self) -> &str {
"Returns the provided value with capability-owned narration."
}
fn parameters_schema(&self) -> serde_json::Value {
json!({
"type": "object",
"properties": {"value": {"type": "string"}},
"required": ["value"],
"additionalProperties": false
})
}
fn narrate(
&self,
_tool_call: &ToolCall,
phase: everruns_core::tool_narration::ToolNarrationPhase,
_locale: Option<&str>,
_ctx: everruns_core::tool_narration::ToolNarrationContext<'_>,
) -> Option<String> {
use everruns_core::tool_narration::ToolNarrationPhase;
match phase {
ToolNarrationPhase::Started | ToolNarrationPhase::Waiting => {
Some("Narrating tool: starting work".to_string())
}
ToolNarrationPhase::Completed => Some("Narrating tool: finished work".to_string()),
ToolNarrationPhase::Failed => Some("Narrating tool: failed".to_string()),
}
}
async fn execute(&self, arguments: serde_json::Value) -> ToolExecutionResult {
ToolExecutionResult::success(json!({
"value": arguments["value"].as_str().unwrap_or_default(),
}))
}
}
struct NarratingCapability;
impl Capability for NarratingCapability {
fn id(&self) -> &str {
"narrating"
}
fn name(&self) -> &str {
"Narrating"
}
fn description(&self) -> &str {
"Test capability whose tool owns its narration."
}
fn status(&self) -> CapabilityStatus {
CapabilityStatus::Available
}
fn tools(&self) -> Vec<Box<dyn Tool>> {
vec![Box::new(NarratingTool)]
}
}
struct ExplicitNarrationHook;
impl everruns_core::capabilities::ToolCallHook for ExplicitNarrationHook {
fn narration(
&self,
_tool_def: Option<&everruns_provider::tool_types::ToolDefinition>,
tool_call: &ToolCall,
_phase: everruns_core::tool_narration::ToolNarrationPhase,
_locale: Option<&str>,
_ctx: everruns_core::tool_narration::ToolNarrationContext<'_>,
) -> Option<String> {
if tool_call.name == "narrating_tool" {
Some("Explicit hook narration wins".to_string())
} else {
None
}
}
}
struct ExplicitNarrationCapability;
impl Capability for ExplicitNarrationCapability {
fn id(&self) -> &str {
"explicit_narration"
}
fn name(&self) -> &str {
"Explicit Narration"
}
fn description(&self) -> &str {
"Test capability with an explicit tool-call narration hook."
}
fn status(&self) -> CapabilityStatus {
CapabilityStatus::Available
}
fn tools(&self) -> Vec<Box<dyn Tool>> {
vec![Box::new(NarratingTool)]
}
fn tool_call_hooks(&self) -> Vec<Arc<dyn everruns_core::capabilities::ToolCallHook>> {
vec![Arc::new(ExplicitNarrationHook)]
}
}
fn harness() -> HarnessDefinition {
HarnessDefinition {
capabilities: vec![everruns_capability::CapabilityRef::new("test_math")],
..HarnessDefinition::new("math", "You are a math harness.")
}
}
fn session(session_id: SessionId, harness_id: HarnessId) -> ExecutionSession {
ExecutionSession {
id: session_id,
workspace_id: everruns_provider::typed_id::WorkspaceId::from_uuid((session_id).uuid()),
organization_id: everruns_core::DEFAULT_ORG_PUBLIC_ID.to_string(),
harness_id,
agent_id: None,
title: Some("Runtime Host".into()),
goal: None,
locale: None,
tags: vec![],
model_id: None,
capabilities: vec![],
tools: vec![],
mcp_servers: Default::default(),
system_prompt: None,
initial_files: vec![],
hints: None,
network_access: None,
max_iterations: None,
parallel_tool_calls: None,
status: SessionExecutionState::Started,
usage: None,
parent_session_id: None,
forked_from_session_id: None,
blueprint_id: None,
blueprint_config: None,
}
}
fn agent(
agent_id: AgentId,
capabilities: Vec<everruns_capability::CapabilityRef>,
) -> AgentDefinition {
AgentDefinition {
display_name: Some("Test Agent".into()),
capabilities,
max_iterations: Some(8),
..AgentDefinition::new(agent_id, "test-agent", "Use tools when needed.")
}
}
fn turn_state(session_id: SessionId, harness_id: HarnessId) -> TurnState {
TurnState {
org_id: 1,
session_id,
harness_id,
agent_id: None,
input_message_id: MessageId::from_uuid(Uuid::now_v7()),
turn_id: Some(TurnId::from_uuid(Uuid::now_v7())),
previous_response_id: None,
iteration: 1,
request_id: None,
started_at: None,
cumulative_usage: None,
tool_call_count: 0,
llm_call_count: 0,
time_to_first_token_ms: None,
final_message_id: None,
final_answer_preview: None,
}
}
async fn advance_from_state(
adapter: &MockHostAdapter,
completed_activity: &str,
state: &TurnState,
output: &serde_json::Value,
pending_user_message_count: usize,
) -> everruns_provider::error::Result<TurnPlan> {
let mut execution = InProcessExecution::new(state.clone());
advance_host_execution(
adapter,
&mut execution,
completed_activity,
output,
pending_user_message_count,
)
.await
}
fn mock_host() -> MockHostAdapter {
let mut capability_registry = CapabilityRegistry::new();
capability_registry.register(TestMathCapability);
let mut driver_registry = DriverRegistry::new();
register_driver(&mut driver_registry);
MockHostAdapter {
capability_registry,
driver_registry,
harness_store: Arc::new(InMemoryHarnessStore::new()),
agent_store: Arc::new(InMemoryAgentStore::new()),
session_store: Arc::new(TestSessionStore::default()),
message_store: Arc::new(InMemoryMessageRetriever::new()),
provider_store: Arc::new(InMemoryProviderStore::new()),
event_emitter: Arc::new(InMemoryEventEmitter::new()),
file_store: Arc::new(InMemorySessionFileStore::new()),
session_task_registry: None,
}
}
async fn set_default_model_spec(adapter: &MockHostAdapter) {
adapter
.provider_store
.set_default_model_spec(ModelSpec::on((DriverId::LlmSim).as_str(), "llmsim-model"))
.await;
}
async fn reason_tool_definitions(
adapter: &MockHostAdapter,
session_id: SessionId,
harness_id: HarnessId,
agent_id: Option<AgentId>,
) -> Vec<everruns_provider::tool_types::ToolDefinition> {
inspect_turn_context(
adapter.harness_store.as_ref(),
adapter.agent_store.as_ref(),
adapter.session_store.as_ref(),
adapter.message_store.as_ref(),
adapter.provider_store.as_ref(),
&adapter.capability_registry,
&adapter.driver_registry,
session_id,
harness_id,
agent_id,
&[],
Some(adapter.file_store.clone()),
)
.await
.expect("reason context")
.runtime_agent
.tools
}
#[tokio::test]
async fn input_activity_emits_lifecycle_events_and_marks_session_active() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let input_message = adapter
.message_store
.add(session_id, InputMessage::user("hello"))
.await
.unwrap();
let turn_id = TurnId::from_uuid(Uuid::now_v7());
execute_input_activity(
&adapter,
1,
InputAtomInput {
context: ExecutionContext::new(session_id, turn_id, input_message.id),
},
)
.await
.unwrap();
let session = adapter
.session_store
.get_session(session_id)
.await
.unwrap()
.unwrap();
assert_eq!(session.status, SessionExecutionState::Active);
let event_types: Vec<_> = adapter
.event_emitter
.events()
.await
.into_iter()
.map(|event| event.data.event_type().to_string())
.collect();
let turn_started_pos = event_types
.iter()
.position(|event_type| event_type == "turn.started")
.expect("turn.started event");
let session_activated_pos = event_types
.iter()
.position(|event_type| event_type == "session.activated")
.expect("session.activated event");
assert!(session_activated_pos < turn_started_pos);
}
#[tokio::test]
async fn act_activity_executes_capability_tools_from_harness_registry() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
let input_message_id = MessageId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let tool_definitions = build_registry(
&adapter.capability_registry,
session_id,
&[everruns_capability::CapabilityRef::new("test_math")],
)
.await
.unwrap()
.tool_definitions();
let result = execute_act_activity(
&adapter,
ActInput {
org_id: Some(1),
context: ExecutionContext::new(
session_id,
TurnId::from_uuid(Uuid::now_v7()),
input_message_id,
),
harness_id,
agent_id: None,
tool_calls: vec![ToolCall {
id: "call_mul".into(),
name: "multiply".into(),
arguments: serde_json::json!({"a": 6, "b": 7}),
}],
tool_definitions,
locale: None,
blueprint_id: None,
network_access: None,
parallel_tool_calls: None,
},
)
.await
.unwrap();
assert_eq!(result.success_count, 1);
assert_eq!(result.error_count, 0);
assert_eq!(result.results.len(), 1);
}
#[cfg(feature = "builtins")]
#[tokio::test]
async fn act_activity_persists_full_output_before_the_hard_limit() {
let mut adapter = mock_host();
adapter
.capability_registry
.register(PersistedOutputCapability);
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
let input_message_id = MessageId::from_uuid(Uuid::now_v7());
let capability = everruns_capability::CapabilityRef::new("persisted_output_fixture");
let mut test_harness = harness();
test_harness.capabilities = vec![capability.clone()];
adapter
.harness_store
.add_harness(harness_id, test_harness)
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let tool_definitions = build_registry(
&adapter.capability_registry,
session_id,
std::slice::from_ref(&capability),
)
.await
.unwrap()
.tool_definitions();
let result = execute_act_activity(
&adapter,
ActInput {
org_id: Some(1),
context: ExecutionContext::new(
session_id,
TurnId::from_uuid(Uuid::now_v7()),
input_message_id,
),
harness_id,
agent_id: None,
tool_calls: vec![ToolCall {
id: "call_persisted_output".into(),
name: "persisted_output_fixture".into(),
arguments: json!({}),
}],
tool_definitions,
locale: None,
blueprint_id: None,
network_access: None,
parallel_tool_calls: None,
},
)
.await
.unwrap();
let payload = result.results[0].result.result.as_ref().unwrap();
assert_eq!(
payload["full_output"],
json!("/workspace/outputs/call_persisted_output.stdout")
);
let persisted = adapter
.file_store
.read_file(session_id, "/outputs/call_persisted_output.stdout")
.await
.unwrap()
.expect("full output persisted before hard limiting");
assert!(persisted.content.unwrap().ends_with("final-marker"));
}
#[tokio::test]
async fn runtime_host_services_reach_final_tool_context_in_parity() {
let mut adapter = mock_host();
adapter
.capability_registry
.register(ContextParityCapability);
adapter.session_task_registry = Some(Arc::new(TestTaskRegistry::default()));
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
let input_message_id = MessageId::from_uuid(Uuid::now_v7());
let mut test_harness = harness();
test_harness.capabilities = vec![everruns_capability::CapabilityRef::new("context_parity")];
adapter
.harness_store
.add_harness(harness_id, test_harness)
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let tool_definitions = build_registry(
&adapter.capability_registry,
session_id,
&[everruns_capability::CapabilityRef::new("context_parity")],
)
.await
.unwrap()
.tool_definitions();
let result = execute_act_activity(
&adapter,
ActInput {
org_id: Some(1),
context: ExecutionContext::new(
session_id,
TurnId::from_uuid(Uuid::now_v7()),
input_message_id,
),
harness_id,
agent_id: None,
tool_calls: vec![ToolCall {
id: "call_context_parity".into(),
name: "context_parity".into(),
arguments: json!({}),
}],
tool_definitions,
locale: None,
blueprint_id: None,
network_access: None,
parallel_tool_calls: None,
},
)
.await
.unwrap();
let payload = result.results[0]
.result
.result
.as_ref()
.expect("context parity payload");
for (service, available) in payload.as_object().expect("object payload") {
assert_eq!(
available,
&json!(true),
"{service} did not reach ToolContext"
);
}
}
#[tokio::test]
async fn act_activity_uses_capability_tool_narration_on_act_path() {
let mut adapter = mock_host();
adapter.capability_registry.register(NarratingCapability);
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
let input_message_id = MessageId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(
harness_id,
HarnessDefinition {
capabilities: vec![everruns_capability::CapabilityRef::new("narrating")],
..harness()
},
)
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let tool_definitions = build_registry(
&adapter.capability_registry,
session_id,
&[everruns_capability::CapabilityRef::new("narrating")],
)
.await
.unwrap()
.tool_definitions();
let result = execute_act_activity(
&adapter,
ActInput {
org_id: Some(1),
context: ExecutionContext::new(
session_id,
TurnId::from_uuid(Uuid::now_v7()),
input_message_id,
),
harness_id,
agent_id: None,
tool_calls: vec![ToolCall {
id: "call_narrate".into(),
name: "narrating_tool".into(),
arguments: json!({"value": "x"}),
}],
tool_definitions,
locale: None,
blueprint_id: None,
network_access: None,
parallel_tool_calls: None,
},
)
.await
.unwrap();
assert_eq!(result.success_count, 1);
let events = adapter.event_emitter.events().await;
let started = events
.iter()
.find_map(|event| match &event.data {
EventData::ToolStarted(data) => Some(data.narration.clone()),
_ => None,
})
.expect("tool.started event");
let completed = events
.iter()
.find_map(|event| match &event.data {
EventData::ToolCompleted(data) => Some(data.narration.clone()),
_ => None,
})
.expect("tool.completed event");
assert_eq!(
started.as_deref(),
Some("Narrating tool: starting work"),
"tool.started must use capability-owned narration, not the generic fallback"
);
assert_eq!(
completed.as_deref(),
Some("Narrating tool: finished work"),
"tool.completed must use capability-owned narration, not the generic fallback"
);
}
#[tokio::test]
async fn act_activity_explicit_tool_call_hook_wins_over_capability_narration() {
let mut adapter = mock_host();
adapter
.capability_registry
.register(ExplicitNarrationCapability);
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
let input_message_id = MessageId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(
harness_id,
HarnessDefinition {
capabilities: vec![everruns_capability::CapabilityRef::new(
"explicit_narration",
)],
..harness()
},
)
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let tool_definitions = build_registry(
&adapter.capability_registry,
session_id,
&[everruns_capability::CapabilityRef::new(
"explicit_narration",
)],
)
.await
.unwrap()
.tool_definitions();
let result = execute_act_activity(
&adapter,
ActInput {
org_id: Some(1),
context: ExecutionContext::new(
session_id,
TurnId::from_uuid(Uuid::now_v7()),
input_message_id,
),
harness_id,
agent_id: None,
tool_calls: vec![ToolCall {
id: "call_narrate".into(),
name: "narrating_tool".into(),
arguments: json!({"value": "x"}),
}],
tool_definitions,
locale: None,
blueprint_id: None,
network_access: None,
parallel_tool_calls: None,
},
)
.await
.unwrap();
assert_eq!(result.success_count, 1);
let events = adapter.event_emitter.events().await;
let started = events
.iter()
.find_map(|event| match &event.data {
EventData::ToolStarted(data) => Some(data.narration.clone()),
_ => None,
})
.expect("tool.started event");
assert_eq!(
started.as_deref(),
Some("Explicit hook narration wins"),
"explicit tool-call hook narration must take precedence over default Tool::narrate()"
);
}
#[tokio::test]
async fn act_activity_agent_session_executes_harness_overlay_tools_from_reason_path() {
let mut adapter = mock_host();
adapter.capability_registry.register(OverlayEchoCapability);
set_default_model_spec(&adapter).await;
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
let agent_id = AgentId::from_uuid(Uuid::now_v7());
let input_message_id = MessageId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(
harness_id,
HarnessDefinition {
capabilities: vec![everruns_capability::CapabilityRef::new("overlay_echo")],
..harness()
},
)
.await;
adapter.agent_store.add_agent(agent(agent_id, vec![])).await;
adapter
.session_store
.insert(ExecutionSession {
agent_id: Some(agent_id),
..session(session_id, harness_id)
})
.await;
let result = execute_act_activity(
&adapter,
ActInput {
org_id: Some(1),
context: ExecutionContext::new(
session_id,
TurnId::from_uuid(Uuid::now_v7()),
input_message_id,
),
harness_id,
agent_id: Some(agent_id),
tool_calls: vec![ToolCall {
id: "call_overlay_echo".into(),
name: "overlay_echo".into(),
arguments: json!({"value": "merged"}),
}],
tool_definitions: reason_tool_definitions(
&adapter,
session_id,
harness_id,
Some(agent_id),
)
.await,
locale: None,
blueprint_id: None,
network_access: None,
parallel_tool_calls: None,
},
)
.await
.unwrap();
assert_eq!(result.success_count, 1);
assert_eq!(result.error_count, 0);
assert_eq!(
result.results[0].result.result.as_ref().unwrap()["value"],
"merged"
);
}
#[tokio::test]
async fn act_activity_agent_session_resolves_transitive_overlay_capabilities() {
let mut adapter = mock_host();
adapter.capability_registry.register(OverlayEchoCapability);
adapter.capability_registry.register(OverlayAliasCapability);
set_default_model_spec(&adapter).await;
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
let agent_id = AgentId::from_uuid(Uuid::now_v7());
let input_message_id = MessageId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(
harness_id,
HarnessDefinition {
capabilities: vec![everruns_capability::CapabilityRef::new("overlay_alias")],
..harness()
},
)
.await;
adapter.agent_store.add_agent(agent(agent_id, vec![])).await;
adapter
.session_store
.insert(ExecutionSession {
agent_id: Some(agent_id),
..session(session_id, harness_id)
})
.await;
let result = execute_act_activity(
&adapter,
ActInput {
org_id: Some(1),
context: ExecutionContext::new(
session_id,
TurnId::from_uuid(Uuid::now_v7()),
input_message_id,
),
harness_id,
agent_id: Some(agent_id),
tool_calls: vec![ToolCall {
id: "call_overlay_dep".into(),
name: "overlay_echo".into(),
arguments: json!({"value": "dependency"}),
}],
tool_definitions: reason_tool_definitions(
&adapter,
session_id,
harness_id,
Some(agent_id),
)
.await,
locale: None,
blueprint_id: None,
network_access: None,
parallel_tool_calls: None,
},
)
.await
.unwrap();
assert_eq!(result.success_count, 1);
assert_eq!(result.error_count, 0);
assert_eq!(
result.results[0].result.result.as_ref().unwrap()["value"],
"dependency"
);
}
#[tokio::test]
async fn act_activity_agent_session_runs_post_tool_hooks_from_merged_overlay() {
let mut adapter = mock_host();
adapter.capability_registry.register(OverlayEchoCapability);
set_default_model_spec(&adapter).await;
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
let agent_id = AgentId::from_uuid(Uuid::now_v7());
let input_message_id = MessageId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(
harness_id,
HarnessDefinition {
capabilities: vec![everruns_capability::CapabilityRef::new("overlay_echo")],
..harness()
},
)
.await;
adapter.agent_store.add_agent(agent(agent_id, vec![])).await;
adapter
.session_store
.insert(ExecutionSession {
agent_id: Some(agent_id),
..session(session_id, harness_id)
})
.await;
let result = execute_act_activity(
&adapter,
ActInput {
org_id: Some(1),
context: ExecutionContext::new(
session_id,
TurnId::from_uuid(Uuid::now_v7()),
input_message_id,
),
harness_id,
agent_id: Some(agent_id),
tool_calls: vec![ToolCall {
id: "call_overlay_hook".into(),
name: "overlay_echo".into(),
arguments: json!({"value": "hook"}),
}],
tool_definitions: reason_tool_definitions(
&adapter,
session_id,
harness_id,
Some(agent_id),
)
.await,
locale: None,
blueprint_id: None,
network_access: None,
parallel_tool_calls: None,
},
)
.await
.unwrap();
assert_eq!(result.success_count, 1);
assert_eq!(result.error_count, 0);
assert_eq!(
result.results[0].result.result.as_ref().unwrap()["hooked"],
true
);
}
#[tokio::test]
async fn act_activity_runs_capability_pre_tool_hook_and_blocks() {
let mut adapter = mock_host();
adapter.capability_registry.register(OverlayEchoCapability);
adapter.capability_registry.register(GateCapability);
set_default_model_spec(&adapter).await;
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
let agent_id = AgentId::from_uuid(Uuid::now_v7());
let input_message_id = MessageId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(
harness_id,
HarnessDefinition {
capabilities: vec![
everruns_capability::CapabilityRef::new("overlay_echo"),
everruns_capability::CapabilityRef::new("gate"),
],
..harness()
},
)
.await;
adapter.agent_store.add_agent(agent(agent_id, vec![])).await;
adapter
.session_store
.insert(ExecutionSession {
agent_id: Some(agent_id),
..session(session_id, harness_id)
})
.await;
let result = execute_act_activity(
&adapter,
ActInput {
org_id: Some(1),
context: ExecutionContext::new(
session_id,
TurnId::from_uuid(Uuid::now_v7()),
input_message_id,
),
harness_id,
agent_id: Some(agent_id),
tool_calls: vec![ToolCall {
id: "call_blocked".into(),
name: "overlay_echo".into(),
arguments: json!({"value": "hook"}),
}],
tool_definitions: reason_tool_definitions(
&adapter,
session_id,
harness_id,
Some(agent_id),
)
.await,
locale: None,
blueprint_id: None,
network_access: None,
parallel_tool_calls: None,
},
)
.await
.unwrap();
assert_eq!(result.success_count, 0);
assert_eq!(result.error_count, 1);
let blocked = &result.results[0].result;
assert!(
blocked.result.is_none(),
"blocked tool should not produce output"
);
let error = blocked
.error
.as_ref()
.expect("blocked tool yields an error");
assert!(
error.contains("denied by test policy"),
"error should carry the hook's reason: {error}"
);
}
#[tokio::test]
async fn act_activity_skips_pre_tool_hooks_from_unavailable_capability() {
let mut adapter = mock_host();
adapter.capability_registry.register(OverlayEchoCapability);
adapter
.capability_registry
.register(ComingSoonGateCapability);
set_default_model_spec(&adapter).await;
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
let agent_id = AgentId::from_uuid(Uuid::now_v7());
let input_message_id = MessageId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(
harness_id,
HarnessDefinition {
capabilities: vec![
everruns_capability::CapabilityRef::new("overlay_echo"),
everruns_capability::CapabilityRef::new("coming_soon_gate"),
],
..harness()
},
)
.await;
adapter.agent_store.add_agent(agent(agent_id, vec![])).await;
adapter
.session_store
.insert(ExecutionSession {
agent_id: Some(agent_id),
..session(session_id, harness_id)
})
.await;
let result = execute_act_activity(
&adapter,
ActInput {
org_id: Some(1),
context: ExecutionContext::new(
session_id,
TurnId::from_uuid(Uuid::now_v7()),
input_message_id,
),
harness_id,
agent_id: Some(agent_id),
tool_calls: vec![ToolCall {
id: "call_unavailable_gate".into(),
name: "overlay_echo".into(),
arguments: json!({"value": "hook"}),
}],
tool_definitions: reason_tool_definitions(
&adapter,
session_id,
harness_id,
Some(agent_id),
)
.await,
locale: None,
blueprint_id: None,
network_access: None,
parallel_tool_calls: None,
},
)
.await
.unwrap();
assert_eq!(result.success_count, 1);
assert_eq!(result.error_count, 0);
assert_eq!(
result.results[0].result.result.as_ref().unwrap()["hooked"],
true
);
}
#[tokio::test]
async fn lifecycle_helper_sets_waiting_for_tool_results_status() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
RuntimeSessionLifecycle::new(adapter.clone(), 1, session_id)
.waiting_for_tool_results()
.await
.unwrap();
let session = adapter
.session_store
.get_session(session_id)
.await
.unwrap()
.unwrap();
assert_eq!(session.status, SessionExecutionState::WaitingForToolResults);
}
#[tokio::test]
async fn execution_propagates_required_lifecycle_write_failures() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
let mut waiting_session = session(session_id, harness_id);
waiting_session.hints = Some(HashMap::from([(
"setup_connection".to_string(),
serde_json::Value::Bool(true),
)]));
adapter.session_store.insert(waiting_session).await;
adapter
.session_store
.fail_status_writes
.store(true, Ordering::SeqCst);
let error = advance_from_state(
&adapter,
"act",
&turn_state(session_id, harness_id),
&serde_json::json!({
"blocked": false,
"waiting_for_tool_results": true
}),
0,
)
.await
.expect_err("required status writes must fail the durable/immediate transition");
assert!(
error
.to_string()
.contains("injected session status failure")
);
}
#[tokio::test]
async fn execution_schedules_reason_after_process_input() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let input = turn_state(session_id, harness_id);
let turn_id = TurnId::from_uuid(Uuid::now_v7());
let output = serde_json::json!({ "turn_id": turn_id.to_string() });
let plan = advance_from_state(&adapter, "process_input", &input, &output, 0)
.await
.unwrap();
match plan {
TurnPlan::ScheduleReason(next) => {
assert_eq!(next.session_id, session_id);
assert_eq!(next.harness_id, harness_id);
assert_eq!(next.turn_id, Some(turn_id));
assert_eq!(next.iteration, 1);
assert_eq!(next.previous_response_id, None);
}
other => panic!("expected ScheduleReason, got {other:?}"),
}
}
#[tokio::test]
async fn execution_schedules_act_after_reason_tool_calls() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let input = turn_state(session_id, harness_id);
let output = serde_json::to_value(ReasonResult {
success: true,
text: "Calling a tool".into(),
tool_calls: vec![ToolCall {
id: "call_mul".into(),
name: "multiply".into(),
arguments: serde_json::json!({ "a": 6, "b": 7 }),
}],
has_tool_calls: true,
tool_definitions: vec![],
max_iterations: 8,
error: None,
user_facing_error: None,
error_disclosure: None,
usage: None,
output_message_id: None,
time_to_first_token_ms: None,
response_id: Some("resp_123".into()),
finish_reason: Some("tool_calls".into()),
locale: Some("en-US".into()),
network_access: None,
parallel_tool_calls: None,
})
.unwrap();
let plan = advance_from_state(&adapter, "reason", &input, &output, 0)
.await
.unwrap();
match plan {
TurnPlan::ScheduleAct(plan) => {
assert_eq!(plan.input.context.session_id, session_id);
assert_eq!(plan.input.context.turn_id, input.turn_id.unwrap());
assert_eq!(plan.input.harness_id, harness_id);
assert_eq!(plan.iteration, 1);
assert_eq!(plan.previous_response_id.as_deref(), Some("resp_123"));
assert_eq!(plan.input.tool_calls.len(), 1);
assert_eq!(plan.input.parallel_tool_calls, None);
}
other => panic!("expected ScheduleAct, got {other:?}"),
}
}
#[tokio::test]
async fn execution_surfaces_max_turn_requests_before_another_act() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let input = turn_state(session_id, harness_id);
let output = serde_json::to_value(ReasonResult {
success: true,
text: "Calling another tool".into(),
tool_calls: vec![ToolCall {
id: "call_mul".into(),
name: "multiply".into(),
arguments: json!({ "a": 6, "b": 7 }),
}],
has_tool_calls: true,
tool_definitions: vec![],
max_iterations: 1,
error: None,
user_facing_error: None,
error_disclosure: None,
usage: None,
output_message_id: None,
time_to_first_token_ms: None,
response_id: Some("resp_limit".into()),
finish_reason: Some("tool_calls".into()),
locale: None,
network_access: None,
parallel_tool_calls: None,
})
.unwrap();
let plan = advance_from_state(&adapter, "reason", &input, &output, 0)
.await
.unwrap();
assert!(matches!(
plan,
TurnPlan::Complete {
stop_reason: TurnStopReason::MaxTurnRequests,
error: None,
}
));
}
#[tokio::test]
async fn execution_threads_parallel_tool_calls_into_act() {
for preference in [Some(true), Some(false), None] {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let input = turn_state(session_id, harness_id);
let output = serde_json::to_value(ReasonResult {
success: true,
text: "Calling a tool".into(),
tool_calls: vec![ToolCall {
id: "call_mul".into(),
name: "multiply".into(),
arguments: serde_json::json!({ "a": 6, "b": 7 }),
}],
has_tool_calls: true,
tool_definitions: vec![],
max_iterations: 8,
error: None,
user_facing_error: None,
error_disclosure: None,
usage: None,
output_message_id: None,
time_to_first_token_ms: None,
response_id: Some("resp_123".into()),
finish_reason: Some("tool_calls".into()),
locale: None,
network_access: None,
parallel_tool_calls: preference,
})
.unwrap();
let plan = advance_from_state(&adapter, "reason", &input, &output, 0)
.await
.unwrap();
match plan {
TurnPlan::ScheduleAct(plan) => {
assert_eq!(
plan.input.parallel_tool_calls, preference,
"parallel_tool_calls must thread through unchanged"
);
}
other => panic!("expected ScheduleAct, got {other:?}"),
}
}
}
#[tokio::test]
async fn execution_schedules_act_with_session_blueprint_id() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(ExecutionSession {
blueprint_id: Some("blueprint.private".into()),
..session(session_id, harness_id)
})
.await;
let input = turn_state(session_id, harness_id);
let output = serde_json::to_value(ReasonResult {
success: true,
text: "Calling a tool".into(),
tool_calls: vec![ToolCall {
id: "call_mul".into(),
name: "multiply".into(),
arguments: serde_json::json!({ "a": 6, "b": 7 }),
}],
has_tool_calls: true,
tool_definitions: vec![],
max_iterations: 8,
error: None,
user_facing_error: None,
error_disclosure: None,
usage: None,
output_message_id: None,
time_to_first_token_ms: None,
response_id: Some("resp_blueprint".into()),
finish_reason: Some("tool_calls".into()),
locale: Some("en-US".into()),
network_access: None,
parallel_tool_calls: None,
})
.unwrap();
let plan = advance_from_state(&adapter, "reason", &input, &output, 0)
.await
.unwrap();
match plan {
TurnPlan::ScheduleAct(plan) => {
assert_eq!(
plan.input.blueprint_id.as_deref(),
Some("blueprint.private")
);
}
other => panic!("expected ScheduleAct, got {other:?}"),
}
}
#[tokio::test]
async fn execution_continues_reason_when_steering_messages_are_pending() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let input = turn_state(session_id, harness_id);
let output = serde_json::to_value(ReasonResult {
success: true,
text: "Continuing".into(),
tool_calls: vec![],
has_tool_calls: false,
tool_definitions: vec![],
max_iterations: 8,
error: None,
user_facing_error: None,
error_disclosure: None,
usage: None,
output_message_id: None,
time_to_first_token_ms: None,
response_id: Some("resp_steer".into()),
finish_reason: Some("stop".into()),
locale: None,
network_access: None,
parallel_tool_calls: None,
})
.unwrap();
let plan = advance_from_state(&adapter, "reason", &input, &output, 2)
.await
.unwrap();
match plan {
TurnPlan::ScheduleReason(next) => {
assert_eq!(next.iteration, 2);
assert_eq!(next.previous_response_id.as_deref(), Some("resp_steer"));
assert_eq!(next.turn_id, input.turn_id);
}
other => panic!("expected ScheduleReason, got {other:?}"),
}
}
#[tokio::test]
async fn execution_emits_turn_completed_summary_fields() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
let final_message_id = MessageId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let mut input = turn_state(session_id, harness_id);
input.started_at = Some(Utc::now());
input.cumulative_usage = Some(TokenUsage::new(10, 4));
input.tool_call_count = 2;
input.llm_call_count = 1;
input.time_to_first_token_ms = Some(35);
let output = serde_json::to_value(ReasonResult {
success: true,
text: "Final answer for export".into(),
tool_calls: vec![],
has_tool_calls: false,
tool_definitions: vec![],
max_iterations: 8,
error: None,
user_facing_error: None,
error_disclosure: None,
usage: Some(TokenUsage::new(20, 8)),
output_message_id: Some(final_message_id),
time_to_first_token_ms: Some(50),
response_id: Some("resp_done".into()),
finish_reason: Some("length".into()),
locale: None,
network_access: None,
parallel_tool_calls: None,
})
.unwrap();
let plan = advance_from_state(&adapter, "reason", &input, &output, 0)
.await
.unwrap();
assert!(matches!(
plan,
TurnPlan::Complete {
stop_reason: everruns_host::TurnStopReason::MaxTokens,
error: None,
}
));
let events = adapter.event_emitter.events().await;
let completed = events
.into_iter()
.find_map(|event| match event.data {
EventData::TurnCompleted(data) => Some(data),
_ => None,
})
.expect("turn.completed event emitted");
assert_eq!(completed.final_message_id, Some(final_message_id));
assert_eq!(
completed.final_answer_preview.as_deref(),
Some("Final answer for export")
);
assert_eq!(completed.time_to_first_token_ms, Some(35));
assert_eq!(completed.tool_call_count, Some(2));
assert_eq!(completed.llm_call_count, Some(2));
assert_eq!(completed.status.as_deref(), Some("completed"));
let usage = completed.usage.expect("aggregated usage");
assert_eq!(usage.input_tokens, 30);
assert_eq!(usage.output_tokens, 12);
assert!(completed.duration_ms.is_some());
}
#[tokio::test]
async fn execution_preserves_reason_failure_message() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let input = turn_state(session_id, harness_id);
let output = serde_json::to_value(ReasonResult {
success: false,
text: "Budget exhausted. 100.00 tokens spent reached the 100.00 tokens limit. Increase the budget to continue."
.into(),
tool_calls: vec![],
has_tool_calls: false,
tool_definitions: vec![],
max_iterations: 8,
error: Some("Budget exhausted".into()),
user_facing_error: None,
error_disclosure: None,
usage: None,
output_message_id: None,
time_to_first_token_ms: None,
response_id: None,
finish_reason: None,
locale: None,
network_access: None,
parallel_tool_calls: None,
})
.unwrap();
let plan = advance_from_state(&adapter, "reason", &input, &output, 0)
.await
.unwrap();
assert!(matches!(plan, TurnPlan::Complete { .. }));
let events = adapter.event_emitter.events().await;
let turn_failed = events
.into_iter()
.find_map(|event| match event.data {
EventData::TurnFailed(data) => Some(data),
_ => None,
})
.expect("turn.failed event emitted");
assert_eq!(
turn_failed.error,
"Budget exhausted. 100.00 tokens spent reached the 100.00 tokens limit. Increase the budget to continue."
);
assert_eq!(
turn_failed.error_code.as_deref(),
Some(user_facing_error_codes::BUDGET_EXHAUSTED)
);
let error_fields = turn_failed.error_fields.expect("budget fields");
assert_eq!(error_fields.get("spent"), Some(&json!(100.0)));
assert_eq!(error_fields.get("limit"), Some(&json!(100.0)));
assert_eq!(error_fields.get("currency"), Some(&json!("tokens")));
}
#[tokio::test]
async fn execution_classifies_missing_api_key_as_provider_misconfigured() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let input = turn_state(session_id, harness_id);
let output = serde_json::to_value(ReasonResult {
success: false,
text: "I encountered an error while processing your request. Please try again later."
.into(),
tool_calls: vec![],
has_tool_calls: false,
tool_definitions: vec![],
max_iterations: 8,
error: Some(
"LLM error: API key is required. Configure the API key in provider settings.".into(),
),
user_facing_error: None,
error_disclosure: None,
usage: None,
output_message_id: None,
time_to_first_token_ms: None,
response_id: None,
finish_reason: None,
locale: None,
network_access: None,
parallel_tool_calls: None,
})
.unwrap();
let plan = advance_from_state(&adapter, "reason", &input, &output, 0)
.await
.unwrap();
assert!(matches!(plan, TurnPlan::Complete { .. }));
let events = adapter.event_emitter.events().await;
let turn_failed = events
.into_iter()
.find_map(|event| match event.data {
EventData::TurnFailed(data) => Some(data),
_ => None,
})
.expect("turn.failed event emitted");
assert_eq!(
turn_failed.error_code.as_deref(),
Some(user_facing_error_codes::PROVIDER_MISCONFIGURED)
);
}
#[tokio::test]
async fn execution_prefers_disclosed_user_facing_error_from_reason() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let input = turn_state(session_id, harness_id);
let output = serde_json::to_value(ReasonResult {
success: false,
text: "I encountered an error while processing your request. Please try again later."
.into(),
tool_calls: vec![],
has_tool_calls: false,
tool_definitions: vec![],
max_iterations: 8,
error: Some("LLM error: OpenAI API error (429): insufficient_quota".into()),
user_facing_error: Some(everruns_provider::user_facing_error::UserFacingError::new(
user_facing_error_codes::PROCESSING_ERROR,
)),
error_disclosure: Some(everruns_provider::user_facing_error::ErrorDisclosure::Generic),
usage: None,
output_message_id: None,
time_to_first_token_ms: None,
response_id: None,
finish_reason: None,
locale: None,
network_access: None,
parallel_tool_calls: None,
})
.unwrap();
let plan = advance_from_state(&adapter, "reason", &input, &output, 0)
.await
.unwrap();
assert!(matches!(plan, TurnPlan::Complete { .. }));
let events = adapter.event_emitter.events().await;
let turn_failed = events
.into_iter()
.find_map(|event| match event.data {
EventData::TurnFailed(data) => Some(data),
_ => None,
})
.expect("turn.failed event emitted");
assert_eq!(
turn_failed.error_code.as_deref(),
Some(user_facing_error_codes::PROCESSING_ERROR)
);
assert_eq!(turn_failed.error_fields, None);
assert_eq!(turn_failed.error_disclosure.as_deref(), Some("generic"));
}
#[tokio::test]
async fn execution_waits_for_tool_results_when_session_hint_requests_it() {
let adapter = mock_host();
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(harness_id, harness())
.await;
let mut host_session = session(session_id, harness_id);
host_session.hints = Some(std::collections::HashMap::from([(
"setup_connection".to_string(),
serde_json::json!(true),
)]));
adapter.session_store.insert(host_session).await;
let input = turn_state(session_id, harness_id);
let output = serde_json::json!({
"waiting_for_tool_results": true,
"blocked": false
});
let plan = advance_from_state(&adapter, "act", &input, &output, 0)
.await
.unwrap();
match plan {
TurnPlan::WaitForToolResults { resume } => {
assert_eq!(resume.iteration, 2);
assert_eq!(resume.turn_id, input.turn_id);
let session = adapter
.session_store
.get_session(session_id)
.await
.unwrap()
.unwrap();
assert_eq!(session.status, SessionExecutionState::WaitingForToolResults);
}
other => panic!("expected WaitForToolResults, got {other:?}"),
}
}
#[cfg(feature = "bashkit")]
struct LifecycleHookCapability {
event: everruns_core::user_hook_types::HookEvent,
command: String,
}
#[cfg(feature = "bashkit")]
impl Capability for LifecycleHookCapability {
fn id(&self) -> &str {
"lifecycle_hook_test"
}
fn name(&self) -> &str {
"Lifecycle Hook Test"
}
fn description(&self) -> &str {
"Test capability contributing a single lifecycle hook."
}
fn status(&self) -> CapabilityStatus {
CapabilityStatus::Available
}
fn user_hooks(&self) -> Vec<everruns_core::user_hook_types::UserHookSpec> {
use everruns_core::user_hook_types::{
ExecutorSpec, HookMatcher, HookSource, OnError, UserHookSpec,
};
vec![UserHookSpec {
id: Some("lc".into()),
event: self.event,
matcher: HookMatcher::default(),
executor: ExecutorSpec::Bash {
command: self.command.clone(),
env: Default::default(),
},
timeout_ms: 5000,
on_error: OnError::Warn,
description: None,
source: HookSource::UserConfig,
}]
}
}
#[cfg(feature = "bashkit")]
#[tokio::test]
async fn user_prompt_submit_hook_blocks_turn_before_reason() {
use everruns_core::user_hook_types::HookEvent;
use everruns_engine::ReasonInput;
use everruns_host::execute_reason_activity;
let mut adapter = mock_host();
adapter.capability_registry.register(LifecycleHookCapability {
event: HookEvent::UserPromptSubmit,
command: r#"echo '{"decision":"block","reason":"policy","user_message":"Your message was blocked."}'"#
.into(),
});
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(
harness_id,
HarnessDefinition {
capabilities: vec![everruns_capability::CapabilityRef::new(
"lifecycle_hook_test",
)],
..harness()
},
)
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let user_message = adapter
.message_store
.add(session_id, InputMessage::user("delete everything"))
.await
.unwrap();
let turn_id = TurnId::from_uuid(Uuid::now_v7());
let result = execute_reason_activity(
&adapter,
1,
ReasonInput {
context: ExecutionContext::new(session_id, turn_id, user_message.id),
harness_id,
agent_id: None,
org_id: 1,
mcp_tool_definitions: vec![],
previous_response_id: None,
iteration: 1,
},
)
.await
.unwrap();
assert!(!result.success, "blocked turn must not succeed");
assert_eq!(result.error.as_deref(), Some("blocked_by_user_prompt_hook"));
assert_eq!(result.text, "Your message was blocked.");
assert!(result.tool_calls.is_empty());
let events = adapter.event_emitter.events().await;
assert!(
events.iter().any(|e| e.event_type == "turn.failed"),
"expected a turn.failed event, got: {:?}",
events.iter().map(|e| &e.event_type).collect::<Vec<_>>()
);
}
#[cfg(feature = "bashkit")]
#[tokio::test]
async fn user_prompt_submit_hook_blocks_injected_mid_turn_message() {
use everruns_core::user_hook_types::HookEvent;
use everruns_engine::ReasonInput;
use everruns_host::execute_reason_activity_with_prompt_messages;
let mut adapter = mock_host();
adapter
.capability_registry
.register(LifecycleHookCapability {
event: HookEvent::UserPromptSubmit,
command: r#"echo '{"decision":"block","reason":"policy"}'"#.into(),
});
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(
harness_id,
HarnessDefinition {
capabilities: vec![everruns_capability::CapabilityRef::new(
"lifecycle_hook_test",
)],
..harness()
},
)
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let original = adapter
.message_store
.add(session_id, InputMessage::user("continue working"))
.await
.unwrap();
let wake = adapter
.message_store
.add(
session_id,
InputMessage::user("Task summary with disallowed content"),
)
.await
.unwrap();
let result = execute_reason_activity_with_prompt_messages(
&adapter,
1,
ReasonInput {
context: ExecutionContext::new(
session_id,
TurnId::from_uuid(Uuid::now_v7()),
original.id,
),
harness_id,
agent_id: None,
org_id: 1,
mcp_tool_definitions: vec![],
previous_response_id: None,
iteration: 3,
},
vec![wake.id],
)
.await
.unwrap();
assert!(!result.success, "injected message must be blocked");
assert_eq!(result.error.as_deref(), Some("blocked_by_user_prompt_hook"));
}
#[cfg(feature = "bashkit")]
#[tokio::test]
async fn user_prompt_submit_hook_allow_does_not_block() {
use everruns_core::user_hook_types::HookEvent;
use everruns_engine::ReasonInput;
use everruns_host::execute_reason_activity;
let mut adapter = mock_host();
set_default_model_spec(&adapter).await;
adapter
.capability_registry
.register(LifecycleHookCapability {
event: HookEvent::UserPromptSubmit,
command: "true".into(), });
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(
harness_id,
HarnessDefinition {
capabilities: vec![everruns_capability::CapabilityRef::new(
"lifecycle_hook_test",
)],
..harness()
},
)
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let user_message = adapter
.message_store
.add(session_id, InputMessage::user("hello"))
.await
.unwrap();
let turn_id = TurnId::from_uuid(Uuid::now_v7());
let result = execute_reason_activity(
&adapter,
1,
ReasonInput {
context: ExecutionContext::new(session_id, turn_id, user_message.id),
harness_id,
agent_id: None,
org_id: 1,
mcp_tool_definitions: vec![],
previous_response_id: None,
iteration: 1,
},
)
.await
.unwrap();
assert_ne!(
result.error.as_deref(),
Some("blocked_by_user_prompt_hook"),
"allow hook must not block the turn"
);
}
#[cfg(feature = "bashkit")]
#[tokio::test]
async fn user_prompt_submit_hook_mutate_rewrites_reason_context() {
use everruns_builtins::InfinityContextCapability;
use everruns_core::user_hook_types::HookEvent;
use everruns_engine::ReasonInput;
use everruns_host::execute_reason_activity;
use everruns_llmsim::{LlmSimConfig, register_driver_with_config};
let mut adapter = mock_host();
let provider_messages = Arc::new(std::sync::Mutex::new(Vec::new()));
register_driver_with_config(
&mut adapter.driver_registry,
LlmSimConfig::echo().with_message_capture(provider_messages.clone()),
);
set_default_model_spec(&adapter).await;
adapter
.capability_registry
.register(LifecycleHookCapability {
event: HookEvent::UserPromptSubmit,
command: r#"echo '{"decision":"mutate","patch":{"message":"sanitized prompt"}}'"#
.into(),
});
adapter
.capability_registry
.register(InfinityContextCapability);
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(
harness_id,
HarnessDefinition {
capabilities: vec![
everruns_capability::CapabilityRef::new("lifecycle_hook_test"),
everruns_capability::CapabilityRef::with_config(
"infinity_context",
json!({ "max_recent_messages": 1 }),
),
],
..harness()
},
)
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
adapter
.message_store
.add(session_id, InputMessage::user("PAST RAW SECRET"))
.await
.unwrap();
let user_message = adapter
.message_store
.add(session_id, InputMessage::user("SECRET prompt"))
.await
.unwrap();
let turn_id = TurnId::from_uuid(Uuid::now_v7());
let result = execute_reason_activity(
&adapter,
1,
ReasonInput {
context: ExecutionContext::new(session_id, turn_id, user_message.id),
harness_id,
agent_id: None,
org_id: 1,
mcp_tool_definitions: vec![],
previous_response_id: None,
iteration: 1,
},
)
.await
.unwrap();
assert!(result.success, "mutated turn should reach the LLM");
assert_eq!(result.text, "Echo: sanitized prompt");
assert!(
result
.tool_definitions
.iter()
.all(|definition| definition.name() != "query_history"),
"raw persisted history must not be exposed when prompt hooks can rewrite user text"
);
let sent = {
let calls = provider_messages.lock().unwrap();
assert!(!calls.is_empty(), "the driver should have been called");
format!("{calls:?}")
};
assert!(
!sent.contains("PAST RAW SECRET"),
"trimmed older history must not reach the provider prompt, got: {sent}"
);
assert!(
sent.contains("sanitized prompt"),
"the current mutated prompt should reach the provider, got: {sent}"
);
let stored_message = adapter
.message_store
.get(session_id, user_message.id)
.await
.unwrap()
.expect("stored user message");
assert_eq!(stored_message.content_to_llm_string(), "SECRET prompt");
}
#[cfg(feature = "bashkit")]
#[tokio::test]
async fn turn_end_hook_fires_on_turn_completion() {
use everruns_core::user_hook_types::HookEvent;
let mut adapter = mock_host();
set_default_model_spec(&adapter).await;
adapter
.capability_registry
.register(LifecycleHookCapability {
event: HookEvent::TurnEnd,
command: "echo turn-ended > /workspace/.turn_end_sentinel".into(),
});
let harness_id = HarnessId::from_uuid(Uuid::now_v7());
let session_id = SessionId::from_uuid(Uuid::now_v7());
adapter
.harness_store
.add_harness(
harness_id,
HarnessDefinition {
capabilities: vec![everruns_capability::CapabilityRef::new(
"lifecycle_hook_test",
)],
..harness()
},
)
.await;
adapter
.session_store
.insert(session(session_id, harness_id))
.await;
let lifecycle = RuntimeSessionLifecycle::new(adapter.clone(), 1, session_id);
lifecycle
.fire_turn_end_hooks(harness_id, None, TurnId::from_uuid(Uuid::now_v7()), true)
.await;
let sentinel = adapter
.file_store
.read_file(session_id, "/.turn_end_sentinel")
.await;
assert!(
sentinel.is_ok() && sentinel.as_ref().unwrap().is_some(),
"turn_end hook should have written the sentinel file, got: {sentinel:?}"
);
}