Skip to main content

vv_agent/runtime/
cycle_runner.rs

1use std::collections::{BTreeMap, BTreeSet};
2
3use serde_json::Value;
4
5use super::context::ExecutionContext;
6use super::hooks::RuntimeHookManager;
7use super::model_calls::{ModelCallCoordinator, ModelCallLedger};
8use super::results::assistant_message_from_response;
9use crate::llm::{LlmClient, LlmError, LlmRequest};
10use crate::memory::{CompactionExhaustedError, MemoryManager};
11use crate::tools::ToolRegistry;
12use crate::types::{AgentTask, CycleRecord, Message, ModelCallOperation};
13
14pub const MAX_PROMPT_TOO_LONG_RETRIES: u32 = 3;
15
16const PROMPT_TOO_LONG_PATTERNS: &[&str] = &[
17    "prompt is too long",
18    "prompt_too_long",
19    "context_length_exceeded",
20    "maximum context length",
21    "request too large",
22    "too many tokens",
23];
24
25pub fn is_prompt_too_long_error(error: &LlmError) -> bool {
26    let text = error.to_string().to_ascii_lowercase();
27    PROMPT_TOO_LONG_PATTERNS
28        .iter()
29        .any(|pattern| text.contains(pattern))
30}
31
32pub struct CycleRunRequest<'a> {
33    pub task: &'a AgentTask,
34    pub messages: Vec<Message>,
35    pub cycle_index: u32,
36    pub memory_manager: &'a mut MemoryManager,
37    pub previous_prompt_tokens: Option<u64>,
38    pub recent_tool_call_ids: Option<&'a BTreeSet<String>>,
39    pub shared_state: Option<&'a BTreeMap<String, Value>>,
40    pub execution_context: Option<&'a ExecutionContext>,
41}
42
43impl<'a> CycleRunRequest<'a> {
44    pub fn new(
45        task: &'a AgentTask,
46        messages: Vec<Message>,
47        cycle_index: u32,
48        memory_manager: &'a mut MemoryManager,
49    ) -> Self {
50        Self {
51            task,
52            messages,
53            cycle_index,
54            memory_manager,
55            previous_prompt_tokens: None,
56            recent_tool_call_ids: None,
57            shared_state: None,
58            execution_context: None,
59        }
60    }
61
62    pub fn with_previous_prompt_tokens(mut self, previous_prompt_tokens: Option<u64>) -> Self {
63        self.previous_prompt_tokens = previous_prompt_tokens;
64        self
65    }
66
67    pub fn with_recent_tool_call_ids(mut self, recent_tool_call_ids: &'a BTreeSet<String>) -> Self {
68        self.recent_tool_call_ids = Some(recent_tool_call_ids);
69        self
70    }
71
72    pub fn with_shared_state(mut self, shared_state: &'a BTreeMap<String, Value>) -> Self {
73        self.shared_state = Some(shared_state);
74        self
75    }
76
77    pub fn with_execution_context(mut self, execution_context: &'a ExecutionContext) -> Self {
78        self.execution_context = Some(execution_context);
79        self
80    }
81}
82
83pub struct CycleRunner<C: LlmClient> {
84    llm_client: C,
85    tool_registry: ToolRegistry,
86    hook_manager: RuntimeHookManager,
87}
88
89impl<C: LlmClient> CycleRunner<C> {
90    pub fn new(llm_client: C, tool_registry: ToolRegistry) -> Self {
91        Self {
92            llm_client,
93            tool_registry,
94            hook_manager: RuntimeHookManager::default(),
95        }
96    }
97
98    pub fn with_hook_manager(mut self, hook_manager: RuntimeHookManager) -> Self {
99        self.hook_manager = hook_manager;
100        self
101    }
102
103    pub fn run_cycle(
104        &self,
105        request: CycleRunRequest<'_>,
106    ) -> Result<(Vec<Message>, CycleRecord), LlmError> {
107        if let Some(context) = request.execution_context {
108            check_context_cancelled(context)?;
109        }
110        let empty_shared_state = BTreeMap::new();
111        let shared_state = request.shared_state.unwrap_or(&empty_shared_state);
112        let pre_compact_messages = self.hook_manager.apply_before_memory_compact(
113            request.task,
114            request.cycle_index,
115            request.messages,
116            shared_state,
117        );
118        let (mut compacted_messages, mut memory_compacted) =
119            request.memory_manager.compact_for_cycle_with_usage(
120                &pre_compact_messages,
121                request.cycle_index,
122                false,
123                request.previous_prompt_tokens,
124                request.recent_tool_call_ids,
125            );
126
127        let mut prompt_too_long_retries = 0;
128        let coordinator = request
129            .execution_context
130            .and_then(|context| context.runtime_state.model_call_coordinator.clone())
131            .unwrap_or_else(|| {
132                let event_handler = request
133                    .execution_context
134                    .and_then(|context| context.event_handler.clone());
135                ModelCallCoordinator::new(
136                    ModelCallLedger::default(),
137                    &request.task.task_id,
138                    &request.task.task_id,
139                    &request.task.task_id,
140                    None,
141                    None,
142                    event_handler,
143                    None,
144                )
145            });
146        let (response, request_messages, request_tool_schemas) = loop {
147            let llm_messages = request
148                .memory_manager
149                .apply_session_memory_context(&compacted_messages);
150            let tool_schemas = self.tool_registry.planned_openai_schemas(request.task);
151            let (request_messages, request_tool_schemas) = self.hook_manager.apply_before_llm(
152                request.task,
153                request.cycle_index,
154                llm_messages,
155                tool_schemas,
156                shared_state,
157            );
158            if let Some(context) = request.execution_context {
159                check_context_cancelled(context)?;
160            }
161            let mut llm_request =
162                LlmRequest::new(request.task.model.clone(), request_messages.clone());
163            llm_request.tools = request_tool_schemas.clone();
164            llm_request.metadata =
165                Value::Object(request.task.metadata.clone().into_iter().collect());
166            llm_request.model_settings = request.task.model_settings.clone();
167            let backend = request
168                .execution_context
169                .and_then(|context| {
170                    context
171                        .metadata
172                        .get("_vv_agent_resolved_backend")
173                        .and_then(Value::as_str)
174                })
175                .unwrap_or("direct");
176            let model = request
177                .execution_context
178                .and_then(|context| {
179                    context
180                        .metadata
181                        .get("_vv_agent_resolved_model")
182                        .and_then(Value::as_str)
183                })
184                .unwrap_or(&request.task.model);
185            let operation_slot = if prompt_too_long_retries == 0 {
186                "main".to_string()
187            } else {
188                format!("prompt_too_long_{prompt_too_long_retries}")
189            };
190            match coordinator.dispatch(
191                ModelCallOperation::AgentCycle,
192                request.cycle_index,
193                &operation_slot,
194                backend,
195                model,
196                &llm_request,
197                || self.llm_client.complete(llm_request.clone()),
198            ) {
199                Ok(dispatch) => break (dispatch.response, request_messages, request_tool_schemas),
200                Err(error) if is_prompt_too_long_error(&error) => {
201                    prompt_too_long_retries += 1;
202                    if prompt_too_long_retries > MAX_PROMPT_TOO_LONG_RETRIES {
203                        return Err(LlmError::CompactionExhausted(
204                            CompactionExhaustedError::new(
205                                prompt_too_long_retries,
206                                Some(error.to_string()),
207                            ),
208                        ));
209                    }
210                    if prompt_too_long_retries == 1 {
211                        (compacted_messages, _) =
212                            request.memory_manager.compact_for_cycle_with_usage(
213                                &compacted_messages,
214                                request.cycle_index,
215                                true,
216                                None,
217                                request.recent_tool_call_ids,
218                            );
219                    } else {
220                        compacted_messages = request.memory_manager.emergency_compact(
221                            &compacted_messages,
222                            (0.2 * f64::from(prompt_too_long_retries)).min(0.95),
223                        );
224                    }
225                    memory_compacted = true;
226                }
227                Err(error) => return Err(error),
228            }
229        };
230
231        if let Some(context) = request.execution_context {
232            check_context_cancelled(context)?;
233        }
234        let response = self.hook_manager.apply_after_llm(
235            request.task,
236            request.cycle_index,
237            &request_messages,
238            &request_tool_schemas,
239            response,
240            shared_state,
241        );
242        let mut next_messages = request_messages;
243        next_messages.push(assistant_message_from_response(&response));
244        let mut cycle = CycleRecord::from_response(request.cycle_index, &response, Vec::new());
245        cycle.memory_compacted = memory_compacted;
246        Ok((next_messages, cycle))
247    }
248}
249
250fn check_context_cancelled(context: &ExecutionContext) -> Result<(), LlmError> {
251    context
252        .check_cancelled()
253        .map_err(|error| LlmError::Request(error.to_string()))
254}