use crate::Controls;
use crate::engine::{
ActOutcome, ActSchedulingFacts, ActivityOutcome, Execution, HostFacts, ReasonResult,
TurnLifecycleEffect, TurnPlan, TurnState, reason_schedules_act,
};
use crate::host::{RuntimeHostAdapter, RuntimeSessionLifecycle};
use chrono::Utc;
use everruns_contracts::error::{AgentLoopError, Result};
use everruns_contracts::typed_id::{MessageId, SessionId, TurnId};
pub async fn advance_host_execution<A: RuntimeHostAdapter, E: Execution>(
adapter: &A,
execution: &mut E,
completed_activity: &str,
output: &serde_json::Value,
pending_user_message_count: usize,
) -> Result<TurnPlan> {
let state = execution.state().clone();
match completed_activity {
"process_input" => {
let turn_id: Option<TurnId> = output
.get("turn_id")
.and_then(|value| value.as_str())
.and_then(|value| value.parse().ok());
let transition = execution.advance(
ActivityOutcome::ProcessInput { turn_id },
pending_user_message_count,
Utc::now(),
HostFacts::default(),
);
perform_effects(adapter, state.org_id, state.session_id, transition.effects).await?;
Ok(transition.plan)
}
"reason" => {
let reason_result: ReasonResult = serde_json::from_value(output.clone())
.map_err(|error| AgentLoopError::Internal(error.into()))?;
let act_scheduling = resolve_act_scheduling(adapter, &state, &reason_result).await?;
let transition = execution.advance(
ActivityOutcome::Reason(Box::new(reason_result)),
pending_user_message_count,
Utc::now(),
HostFacts {
act_scheduling,
..HostFacts::default()
},
);
perform_effects(adapter, state.org_id, state.session_id, transition.effects).await?;
Ok(transition.plan)
}
"act" => {
let client_tool_calls: Vec<everruns_contracts::tool_types::ToolCall> = output
.get("client_tool_calls")
.cloned()
.map(serde_json::from_value)
.transpose()
.map_err(|error| AgentLoopError::Internal(error.into()))?
.unwrap_or_default();
let ask_user_calls = pending_ask_user_calls_from_slice(&client_tool_calls);
let outcome = ActOutcome {
blocked: output
.get("blocked")
.and_then(|value| value.as_bool())
.unwrap_or(false),
waiting_for_tool_results: output
.get("waiting_for_tool_results")
.and_then(|value| value.as_bool())
.unwrap_or(false),
waiting_for_url_elicitation: output
.get("waiting_for_url_elicitation")
.and_then(|value| value.as_bool())
.unwrap_or(false),
waiting_for_ask_user: !ask_user_calls.is_empty(),
waiting_for_tool_approval: crate::engine::has_pending_tool_approval(
&client_tool_calls,
),
};
let hints = resolve_pause_hints(
adapter,
state.org_id,
state.session_id,
state.input_message_id,
outcome,
)
.await?;
let transition = execution.advance(
ActivityOutcome::Act(outcome),
pending_user_message_count,
Utc::now(),
HostFacts {
setup_connection_hint_enabled: hints.setup_connection,
url_elicitation_hint_enabled: hints.url_elicitation,
ask_user_hint_enabled: hints.ask_user,
ask_user_calls,
..HostFacts::default()
},
);
perform_effects(adapter, state.org_id, state.session_id, transition.effects).await?;
Ok(transition.plan)
}
other => Err(AgentLoopError::config(format!(
"Unknown activity type completed: {other}"
))),
}
}
pub(crate) async fn resolve_act_scheduling<A: RuntimeHostAdapter>(
adapter: &A,
state: &TurnState,
reason_result: &ReasonResult,
) -> Result<Option<ActSchedulingFacts>> {
if !reason_schedules_act(state, reason_result) {
return Ok(None);
}
let session = adapter
.session_store(state.org_id)
.get_session(state.session_id)
.await?;
Ok(Some(ActSchedulingFacts {
blueprint_id: session.as_ref().and_then(|s| s.blueprint_id.clone()),
workspace_id: session.as_ref().map(|s| s.workspace_id),
}))
}
#[derive(Debug, Clone, Copy, Default)]
pub(crate) struct PauseHints {
pub setup_connection: bool,
pub url_elicitation: bool,
pub ask_user: bool,
}
pub(crate) async fn resolve_pause_hints<A: RuntimeHostAdapter>(
adapter: &A,
org_id: i64,
session_id: SessionId,
input_message_id: MessageId,
outcome: ActOutcome,
) -> Result<PauseHints> {
if outcome.blocked || !outcome.waiting_for_tool_results {
return Ok(PauseHints::default());
}
match adapter.session_store(org_id).get_session(session_id).await {
Ok(Some(session)) => {
let message = adapter
.message_store()
.get(session_id, input_message_id)
.await?;
let message_hints = message
.as_ref()
.and_then(|message| message.controls.as_ref())
.and_then(|controls| controls.hints.as_ref());
let hints = Controls::resolve_hints(session.hints.as_ref(), message_hints);
let flag = |name: &str| {
hints
.get(name)
.and_then(|value| value.as_bool())
.unwrap_or(false)
};
Ok(PauseHints {
setup_connection: flag("setup_connection"),
url_elicitation: flag("url_elicitation"),
ask_user: flag("ask_user"),
})
}
_ => Ok(PauseHints::default()),
}
}
pub(crate) async fn perform_effects<A: RuntimeHostAdapter>(
adapter: &A,
org_id: i64,
session_id: SessionId,
effects: Vec<TurnLifecycleEffect>,
) -> Result<()> {
if effects.is_empty() {
return Ok(());
}
let lifecycle = RuntimeSessionLifecycle::new(adapter.clone(), org_id, session_id);
for effect in effects {
match effect {
TurnLifecycleEffect::TurnCompleted {
input_message_id,
data,
} => {
lifecycle
.emit_turn_completed(input_message_id, data)
.await?;
}
TurnLifecycleEffect::SessionIdled {
turn_id,
input_message_id,
iterations,
usage,
} => {
lifecycle
.emit_session_idled(turn_id, input_message_id, iterations, usage)
.await?;
}
TurnLifecycleEffect::TurnFailedWithDisclosure {
turn_id,
input_message_id,
text,
user_error,
disclosure,
} => {
lifecycle
.turn_failed_with_disclosure(
turn_id,
input_message_id,
&text,
user_error.as_ref(),
disclosure,
)
.await?;
}
TurnLifecycleEffect::FireTurnEndHooks {
harness_id,
agent_id,
turn_id,
success,
} => {
lifecycle
.fire_turn_end_hooks(harness_id, agent_id, turn_id, success)
.await;
}
TurnLifecycleEffect::WaitingForToolResults => {
lifecycle.waiting_for_tool_results().await?;
}
TurnLifecycleEffect::ResolveAskUserUnattended {
turn_id,
input_message_id,
calls,
} => {
lifecycle
.resolve_ask_user_unattended(turn_id, input_message_id, calls)
.await?;
}
}
}
Ok(())
}
pub(crate) fn pending_ask_user_calls(
act_result: &crate::engine::ActResult,
) -> Vec<(String, serde_json::Value)> {
pending_ask_user_calls_from_slice(&act_result.client_tool_calls)
}
pub(crate) fn pending_ask_user_calls_from_slice(
client_tool_calls: &[everruns_contracts::tool_types::ToolCall],
) -> Vec<(String, serde_json::Value)> {
client_tool_calls
.iter()
.filter(|call| call.name == everruns_contracts::ASK_USER_TOOL_NAME)
.map(|call| (call.id.clone(), call.arguments.clone()))
.collect()
}
pub(crate) fn has_pending_ask_user(act_result: &crate::engine::ActResult) -> bool {
act_result
.client_tool_calls
.iter()
.any(|call| call.name == everruns_contracts::ASK_USER_TOOL_NAME)
}
pub(crate) fn act_outcome(act_result: &crate::engine::ActResult) -> ActOutcome {
ActOutcome {
blocked: act_result.blocked,
waiting_for_tool_results: act_result.waiting_for_tool_results,
waiting_for_url_elicitation: act_result.waiting_for_url_elicitation,
waiting_for_ask_user: has_pending_ask_user(act_result),
waiting_for_tool_approval: crate::engine::has_pending_tool_approval(
&act_result.client_tool_calls,
),
}
}