pub(crate) mod actions;
mod checkpoint;
#[cfg(feature = "background")]
mod compaction;
mod inference;
mod logical_inference;
mod orchestrator;
#[cfg(feature = "parallel-tools")]
pub mod parallel_merge;
mod resume;
mod setup;
mod step;
mod stream_policy;
#[cfg(test)]
mod tests;
use std::sync::Arc;
use crate::cancellation::CancellationToken;
use crate::checkpoint_store::RuntimeCheckpointStore;
use crate::phase::{ExecutionEnv, PhaseRuntime};
use crate::registry::AgentResolver;
use crate::state::MutationBatch;
use async_trait::async_trait;
use awaken_runtime_contract::StateError;
use awaken_runtime_contract::contract::event_sink::EventSink;
use awaken_runtime_contract::contract::identity::RunIdentity;
use awaken_runtime_contract::contract::inference::InferenceOverride;
use awaken_runtime_contract::contract::message::{DeliveryBoundary, Message};
use awaken_runtime_contract::contract::suspension::ToolCallResume;
use awaken_runtime_contract::contract::tool::{ToolResult, ToolStatus};
use futures::channel::mpsc;
use serde_json::Value;
use crate::agent::state::{RunLifecycle, ToolCallStates};
pub use actions::LoopActionHandlersPlugin;
pub use checkpoint::CommitWiring;
pub(crate) use checkpoint::{CommitAppendError, commit_checkpoint_appending};
pub use resume::prepare_resume;
pub struct LoopStatePlugin;
impl crate::plugins::Plugin for LoopStatePlugin {
fn descriptor(&self) -> crate::plugins::PluginDescriptor {
crate::plugins::PluginDescriptor {
name: "__loop_state",
}
}
fn register(
&self,
r: &mut crate::plugins::PluginRegistrar,
) -> Result<(), awaken_runtime_contract::StateError> {
use crate::agent::state::{ContextMessageStore, ContextThrottleState};
use crate::state::{KeyScope, StateKeyOptions};
r.register_key::<RunLifecycle>(StateKeyOptions::default())?;
r.register_key::<ToolCallStates>(StateKeyOptions {
scope: KeyScope::Thread,
persistent: true,
..StateKeyOptions::default()
})?;
r.register_key::<ContextThrottleState>(StateKeyOptions::default())?;
r.register_key::<ContextMessageStore>(StateKeyOptions::default())?;
r.register_key::<crate::agent::state::PendingWorkKey>(StateKeyOptions::default())?;
Ok(())
}
}
#[derive(Debug, thiserror::Error)]
pub enum AgentLoopError {
#[error("inference failed: {0}")]
InferenceFailed(String),
#[error("inference failed: {0}")]
Inference(#[from] awaken_runtime_contract::contract::executor::InferenceExecutionError),
#[error("storage failed: {0}")]
StorageError(String),
#[error("phase error: {0}")]
PhaseError(#[from] awaken_runtime_contract::StateError),
#[error("runtime error: {0}")]
RuntimeError(#[from] crate::error::RuntimeError),
#[error("invalid activation: {0}")]
InvalidActivation(String),
#[error("invalid resume: {0}")]
InvalidResume(String),
}
impl From<crate::execution::executor::ToolExecutorError> for AgentLoopError {
fn from(e: crate::execution::executor::ToolExecutorError) -> Self {
Self::InferenceFailed(e.to_string())
}
}
#[derive(Debug)]
pub struct AgentRunResult {
pub run_id: String,
pub response: String,
pub termination: awaken_runtime_contract::contract::lifecycle::TerminationReason,
pub steps: usize,
}
#[derive(Debug, Clone, Default)]
pub struct PendingBoundaryFreeze {
pub messages: Vec<Message>,
}
#[async_trait]
pub trait PendingBoundaryHandler: Send + Sync {
async fn stage_pending_messages(
&self,
boundary: DeliveryBoundary,
messages: Vec<Message>,
) -> Result<(), AgentLoopError>;
async fn freeze_pending_boundary(
&self,
boundary: DeliveryBoundary,
) -> Result<Option<PendingBoundaryFreeze>, AgentLoopError>;
}
pub(crate) use awaken_runtime_contract::now_ms;
fn commit_update<S: crate::state::StateKey>(
store: &crate::state::StateStore,
update: S::Update,
) -> Result<(), awaken_runtime_contract::StateError> {
let mut patch = MutationBatch::new();
patch.update::<S>(update);
store.commit(patch)?;
clear_pending_scheduled_actions_for_terminal_run::<S>(store)?;
Ok(())
}
fn clear_pending_scheduled_actions_for_terminal_run<S: crate::state::StateKey>(
store: &crate::state::StateStore,
) -> Result<(), awaken_runtime_contract::StateError> {
if S::KEY != "__runtime.run_lifecycle" {
return Ok(());
}
let Some(lifecycle) = store.read::<RunLifecycle>() else {
return Ok(());
};
if !lifecycle.status.is_terminal() {
return Ok(());
}
let Some(pending) = store.read::<awaken_runtime_contract::model::PendingScheduledActions>()
else {
return Ok(());
};
if pending.is_empty() {
return Ok(());
}
let mut cleanup = MutationBatch::new();
for action in pending {
cleanup.update::<awaken_runtime_contract::model::PendingScheduledActions>(
awaken_runtime_contract::model::ScheduledActionQueueUpdate::Remove { id: action.id },
);
}
store.commit(cleanup)?;
Ok(())
}
fn tool_result_to_content(result: &ToolResult) -> String {
match &result.message {
Some(msg) => msg.clone(),
None => serde_json::to_string(&result.data).unwrap_or_default(),
}
}
fn tool_result_to_resume_payload(result: &ToolResult) -> Value {
match result.status {
ToolStatus::Success => {
if result.metadata.is_empty() {
result.data.clone()
} else {
serde_json::json!({
"data": result.data,
"metadata": result.metadata,
})
}
}
ToolStatus::Error => {
if let Some(message) = result.message.as_ref() {
serde_json::json!({ "error": message })
} else {
result.data.clone()
}
}
ToolStatus::Pending => Value::Null,
}
}
pub struct AgentLoopParams<'a> {
pub resolver: &'a dyn AgentResolver,
pub agent_id: &'a str,
pub runtime: &'a PhaseRuntime,
pub sink: Arc<dyn EventSink>,
pub checkpoint_store: Option<&'a dyn RuntimeCheckpointStore>,
pub commit: checkpoint::CommitWiring<'a>,
pub messages: Vec<Message>,
pub run_identity: RunIdentity,
pub cancellation_token: Option<CancellationToken>,
pub decision_rx: Option<mpsc::UnboundedReceiver<Vec<(String, ToolCallResume)>>>,
pub overrides: Option<InferenceOverride>,
pub frontend_tools: Vec<awaken_runtime_contract::contract::tool::ToolDescriptor>,
pub inbox: Option<crate::inbox::InboxReceiver>,
pub is_continuation: bool,
pub initial_state_seed: Option<awaken_runtime_contract::state::PersistedState>,
}
pub fn build_agent_env(
plugins: &[Arc<dyn crate::plugins::Plugin>],
agent: &crate::registry::ResolvedAgent,
) -> Result<ExecutionEnv, StateError> {
let stop_policies = crate::policies::policies_from_specs(agent.stop_conditions());
let mut all_plugins = crate::registry::resolve::inject_default_plugins_with_stop_policies(
plugins.to_vec(),
agent.max_rounds(),
stop_policies,
);
if let Some(policy) = agent.context_policy() {
let transform_config = agent
.spec
.config::<crate::context::ContextTransformConfigKey>()
.unwrap_or_default();
all_plugins.push(Arc::new(
crate::context::ContextTransformPlugin::with_config(policy.clone(), transform_config),
));
}
ExecutionEnv::from_plugins(&all_plugins, &std::collections::HashSet::new())
}
pub async fn run_agent_loop(params: AgentLoopParams<'_>) -> Result<AgentRunResult, AgentLoopError> {
orchestrator::run_agent_loop_impl(params, None, None).await
}
pub(crate) async fn run_agent_loop_with_pending_boundary(
params: AgentLoopParams<'_>,
thread_ctx: Option<crate::ThreadContextSnapshot>,
pending_boundary: Option<Arc<dyn PendingBoundaryHandler>>,
) -> Result<AgentRunResult, AgentLoopError> {
orchestrator::run_agent_loop_impl(params, thread_ctx, pending_boundary).await
}