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}