awaken_runtime/loop_runner/
mod.rs1pub(crate) mod actions;
7mod checkpoint;
8#[cfg(feature = "background")]
9mod compaction;
10mod inference;
11mod logical_inference;
12mod orchestrator;
13#[cfg(feature = "parallel-tools")]
14pub mod parallel_merge;
15mod resume;
16mod setup;
17mod step;
18mod stream_policy;
19
20#[cfg(test)]
21mod tests;
22
23use std::sync::Arc;
24
25use crate::cancellation::CancellationToken;
26use crate::checkpoint_store::RuntimeCheckpointStore;
27use crate::phase::{ExecutionEnv, PhaseRuntime};
28use crate::registry::AgentResolver;
29use crate::state::MutationBatch;
30use async_trait::async_trait;
31use awaken_runtime_contract::StateError;
32use awaken_runtime_contract::contract::event_sink::EventSink;
33use awaken_runtime_contract::contract::identity::RunIdentity;
34use awaken_runtime_contract::contract::inference::InferenceOverride;
35use awaken_runtime_contract::contract::message::{DeliveryBoundary, Message};
36use awaken_runtime_contract::contract::suspension::ToolCallResume;
37use awaken_runtime_contract::contract::tool::{ToolResult, ToolStatus};
38use futures::channel::mpsc;
39use serde_json::Value;
40
41use crate::agent::state::{RunLifecycle, ToolCallStates};
42
43pub use actions::LoopActionHandlersPlugin;
45pub use checkpoint::CommitWiring;
46pub(crate) use checkpoint::{CommitAppendError, commit_checkpoint_appending};
47pub use resume::prepare_resume;
48
49pub struct LoopStatePlugin;
53
54impl crate::plugins::Plugin for LoopStatePlugin {
55 fn descriptor(&self) -> crate::plugins::PluginDescriptor {
56 crate::plugins::PluginDescriptor {
57 name: "__loop_state",
58 }
59 }
60
61 fn register(
62 &self,
63 r: &mut crate::plugins::PluginRegistrar,
64 ) -> Result<(), awaken_runtime_contract::StateError> {
65 use crate::agent::state::{ContextMessageStore, ContextThrottleState};
66 use crate::state::{KeyScope, StateKeyOptions};
67
68 r.register_key::<RunLifecycle>(StateKeyOptions::default())?;
69 r.register_key::<ToolCallStates>(StateKeyOptions {
70 scope: KeyScope::Thread,
71 persistent: true,
72 ..StateKeyOptions::default()
73 })?;
74 r.register_key::<ContextThrottleState>(StateKeyOptions::default())?;
75 r.register_key::<ContextMessageStore>(StateKeyOptions::default())?;
76 r.register_key::<crate::agent::state::PendingWorkKey>(StateKeyOptions::default())?;
77
78 Ok(())
79 }
80}
81
82#[derive(Debug, thiserror::Error)]
84pub enum AgentLoopError {
85 #[error("inference failed: {0}")]
86 InferenceFailed(String),
87 #[error("inference failed: {0}")]
97 Inference(#[from] awaken_runtime_contract::contract::executor::InferenceExecutionError),
98 #[error("storage failed: {0}")]
99 StorageError(String),
100 #[error("phase error: {0}")]
101 PhaseError(#[from] awaken_runtime_contract::StateError),
102 #[error("runtime error: {0}")]
103 RuntimeError(#[from] crate::error::RuntimeError),
104 #[error("invalid activation: {0}")]
105 InvalidActivation(String),
106 #[error("invalid resume: {0}")]
107 InvalidResume(String),
108}
109
110impl From<crate::execution::executor::ToolExecutorError> for AgentLoopError {
111 fn from(e: crate::execution::executor::ToolExecutorError) -> Self {
112 Self::InferenceFailed(e.to_string())
113 }
114}
115
116#[derive(Debug)]
118pub struct AgentRunResult {
119 pub run_id: String,
120 pub response: String,
121 pub termination: awaken_runtime_contract::contract::lifecycle::TerminationReason,
122 pub steps: usize,
123}
124
125#[derive(Debug, Clone, Default)]
127pub struct PendingBoundaryFreeze {
128 pub messages: Vec<Message>,
129}
130
131#[async_trait]
137pub trait PendingBoundaryHandler: Send + Sync {
138 async fn stage_pending_messages(
139 &self,
140 boundary: DeliveryBoundary,
141 messages: Vec<Message>,
142 ) -> Result<(), AgentLoopError>;
143
144 async fn freeze_pending_boundary(
145 &self,
146 boundary: DeliveryBoundary,
147 ) -> Result<Option<PendingBoundaryFreeze>, AgentLoopError>;
148}
149
150pub(crate) use awaken_runtime_contract::now_ms;
153
154fn commit_update<S: crate::state::StateKey>(
155 store: &crate::state::StateStore,
156 update: S::Update,
157) -> Result<(), awaken_runtime_contract::StateError> {
158 let mut patch = MutationBatch::new();
159 patch.update::<S>(update);
160 store.commit(patch)?;
161 clear_pending_scheduled_actions_for_terminal_run::<S>(store)?;
162 Ok(())
163}
164
165fn clear_pending_scheduled_actions_for_terminal_run<S: crate::state::StateKey>(
166 store: &crate::state::StateStore,
167) -> Result<(), awaken_runtime_contract::StateError> {
168 if S::KEY != "__runtime.run_lifecycle" {
169 return Ok(());
170 }
171 let Some(lifecycle) = store.read::<RunLifecycle>() else {
172 return Ok(());
173 };
174 if !lifecycle.status.is_terminal() {
175 return Ok(());
176 }
177 let Some(pending) = store.read::<awaken_runtime_contract::model::PendingScheduledActions>()
178 else {
179 return Ok(());
180 };
181 if pending.is_empty() {
182 return Ok(());
183 }
184
185 let mut cleanup = MutationBatch::new();
186 for action in pending {
187 cleanup.update::<awaken_runtime_contract::model::PendingScheduledActions>(
188 awaken_runtime_contract::model::ScheduledActionQueueUpdate::Remove { id: action.id },
189 );
190 }
191 store.commit(cleanup)?;
192 Ok(())
193}
194
195fn tool_result_to_content(result: &ToolResult) -> String {
196 match &result.message {
197 Some(msg) => msg.clone(),
198 None => serde_json::to_string(&result.data).unwrap_or_default(),
199 }
200}
201
202fn tool_result_to_resume_payload(result: &ToolResult) -> Value {
203 match result.status {
204 ToolStatus::Success => {
205 if result.metadata.is_empty() {
206 result.data.clone()
207 } else {
208 serde_json::json!({
209 "data": result.data,
210 "metadata": result.metadata,
211 })
212 }
213 }
214 ToolStatus::Error => {
215 if let Some(message) = result.message.as_ref() {
216 serde_json::json!({ "error": message })
217 } else {
218 result.data.clone()
219 }
220 }
221 ToolStatus::Pending => Value::Null,
222 }
223}
224
225pub struct AgentLoopParams<'a> {
227 pub resolver: &'a dyn AgentResolver,
229 pub agent_id: &'a str,
231 pub runtime: &'a PhaseRuntime,
233 pub sink: Arc<dyn EventSink>,
235 pub checkpoint_store: Option<&'a dyn RuntimeCheckpointStore>,
237 pub commit: checkpoint::CommitWiring<'a>,
239 pub messages: Vec<Message>,
241 pub run_identity: RunIdentity,
243 pub cancellation_token: Option<CancellationToken>,
245 pub decision_rx: Option<mpsc::UnboundedReceiver<Vec<(String, ToolCallResume)>>>,
247 pub overrides: Option<InferenceOverride>,
249 pub frontend_tools: Vec<awaken_runtime_contract::contract::tool::ToolDescriptor>,
255 pub inbox: Option<crate::inbox::InboxReceiver>,
257 pub is_continuation: bool,
260 pub initial_state_seed: Option<awaken_runtime_contract::state::PersistedState>,
266}
267
268pub fn build_agent_env(
276 plugins: &[Arc<dyn crate::plugins::Plugin>],
277 agent: &crate::registry::ResolvedAgent,
278) -> Result<ExecutionEnv, StateError> {
279 let stop_policies = crate::policies::policies_from_specs(agent.stop_conditions());
280 let mut all_plugins = crate::registry::resolve::inject_default_plugins_with_stop_policies(
281 plugins.to_vec(),
282 agent.max_rounds(),
283 stop_policies,
284 );
285
286 if let Some(policy) = agent.context_policy() {
287 let transform_config = agent
288 .spec
289 .config::<crate::context::ContextTransformConfigKey>()
290 .unwrap_or_default();
291 all_plugins.push(Arc::new(
292 crate::context::ContextTransformPlugin::with_config(policy.clone(), transform_config),
293 ));
294 }
295
296 ExecutionEnv::from_plugins(&all_plugins, &std::collections::HashSet::new())
297}
298
299pub async fn run_agent_loop(params: AgentLoopParams<'_>) -> Result<AgentRunResult, AgentLoopError> {
305 orchestrator::run_agent_loop_impl(params, None, None).await
306}
307
308pub(crate) async fn run_agent_loop_with_pending_boundary(
309 params: AgentLoopParams<'_>,
310 thread_ctx: Option<crate::ThreadContextSnapshot>,
311 pending_boundary: Option<Arc<dyn PendingBoundaryHandler>>,
312) -> Result<AgentRunResult, AgentLoopError> {
313 orchestrator::run_agent_loop_impl(params, thread_ctx, pending_boundary).await
314}