1use crate::exec::async_command::{AsyncProcessRunner, ProcessOptions, StreamCaptureConfig};
32use crate::exec::sdk_ipc::{ToolIpcHandler, ToolResponse};
33use crate::mcp::McpToolExecutor;
34use crate::utils::async_utils;
35use crate::utils::file_utils::{ensure_dir_exists, write_file_with_context};
36use anyhow::{Context, Result};
37use async_trait::async_trait;
38use hashbrown::HashMap;
39use serde_json::Value;
40use std::ffi::OsString;
41use std::path::PathBuf;
42use std::sync::Arc;
43use std::time::{Duration, Instant};
44use tokio::task::JoinHandle;
45use tracing::{debug, info};
46
47#[derive(Debug, Clone, Copy, PartialEq, Eq)]
49pub enum Language {
50 Python3,
52 JavaScript,
54}
55
56impl Language {
57 pub fn as_str(&self) -> &'static str {
59 match self {
60 Self::Python3 => "python3",
61 Self::JavaScript => "javascript",
62 }
63 }
64
65 pub fn interpreter(&self) -> &'static str {
67 match self {
68 Self::Python3 => "python3",
69 Self::JavaScript => "node",
70 }
71 }
72
73 pub fn detect_python_interpreter(workspace_root: &std::path::Path) -> String {
75 if let Ok(venv_python) = std::env::var("VIRTUAL_ENV") {
77 let venv_bin = PathBuf::from(venv_python).join("bin").join("python");
78 if venv_bin.exists() {
79 debug!("Using venv Python: {:?}", venv_bin);
80 return venv_bin.to_string_lossy().into_owned();
81 }
82 }
83
84 let workspace_venv = workspace_root.join(".venv").join("bin").join("python");
86 if workspace_venv.exists() {
87 debug!("Using workspace .venv Python: {:?}", workspace_venv);
88 return workspace_venv.to_string_lossy().into_owned();
89 }
90
91 if let Ok(system_python) = which::which("python3") {
93 debug!("Using system python3: {:?}", system_python);
94 return system_python.to_string_lossy().into_owned();
95 }
96
97 if which::which("uv").is_ok() {
99 debug!("Using uv for Python execution");
100 return "uv".to_string();
101 }
102
103 debug!("Using system python3");
106 "python3".to_string()
107 }
108}
109
110#[derive(Debug, Clone)]
114pub struct BuiltinToolInfo {
115 pub name: String,
117 pub description: String,
119}
120
121#[async_trait]
123pub trait BuiltinToolExecutor: Send + Sync {
124 async fn execute_builtin_tool(&self, tool_name: &str, args: &Value) -> Result<Value>;
126 fn list_builtin_tools(&self) -> Result<Vec<BuiltinToolInfo>>;
128}
129
130#[derive(Debug, Clone, serde::Serialize)]
132pub struct ExecutionResult {
133 pub exit_code: i32,
135 pub stdout: String,
137 pub stderr: String,
139 pub json_result: Option<Value>,
141 pub duration_ms: u128,
143}
144
145#[derive(Debug, Clone)]
147pub struct ExecutionConfig {
148 pub timeout_secs: u64,
150 pub max_output_bytes: usize,
152}
153
154impl Default for ExecutionConfig {
155 fn default() -> Self {
156 Self {
157 timeout_secs: 30,
158 max_output_bytes: 10 * 1024 * 1024, }
160 }
161}
162
163pub struct CodeExecutor {
165 language: Language,
166 mcp_client: Arc<dyn McpToolExecutor>,
167 builtin_executor: Option<Arc<dyn BuiltinToolExecutor>>,
168 config: ExecutionConfig,
169 workspace_root: PathBuf,
170 enable_pii_protection: bool,
171}
172
173impl CodeExecutor {
174 pub fn new(language: Language, mcp_client: Arc<dyn McpToolExecutor>, workspace_root: PathBuf) -> Self {
176 Self {
177 language,
178 mcp_client,
179 builtin_executor: None,
180 config: ExecutionConfig::default(),
181 workspace_root,
182 enable_pii_protection: false,
183 }
184 }
185
186 pub fn with_builtin_executor(mut self, executor: Option<Arc<dyn BuiltinToolExecutor>>) -> Self {
190 self.builtin_executor = executor;
191 self
192 }
193
194 pub async fn execute(&self, code: &str) -> Result<ExecutionResult> {
205 info!(language = self.language.as_str(), timeout_secs = self.config.timeout_secs, "Executing code snippet");
206
207 let start = Instant::now();
208
209 let ipc_dir = self.workspace_root.join(".vtcode").join("ipc");
211 ensure_dir_exists(&ipc_dir).await?;
212
213 let sdk = self.generate_sdk().await.context("failed to generate SDK")?;
215
216 let complete_code = match self.language {
218 Language::Python3 => self.prepare_python_code(&sdk, code)?,
219 Language::JavaScript => self.prepare_javascript_code(&sdk, code)?,
220 };
221
222 let code_temp_dir = self.workspace_root.join(".vtcode").join("code_temp");
225 ensure_dir_exists(&code_temp_dir).await?;
226
227 let timestamp = std::time::SystemTime::now()
229 .duration_since(std::time::UNIX_EPOCH)
230 .unwrap_or_default()
231 .as_micros();
232 let ext = match self.language {
233 Language::Python3 => "py",
234 Language::JavaScript => "js",
235 };
236 let code_file = code_temp_dir.join(format!("exec_{timestamp}.{ext}"));
237
238 write_file_with_context(&code_file, &complete_code, "temporary code file").await?;
239
240 debug!(
241 language = self.language.as_str(),
242 code_file = ?code_file,
243 "Wrote code to temporary file"
244 );
245
246 let mut env = HashMap::new();
248
249 env.insert(OsString::from("VTCODE_IPC_DIR"), OsString::from(&*ipc_dir.to_string_lossy()));
251
252 let mut ipc_handler = ToolIpcHandler::new(ipc_dir.clone());
254 if self.enable_pii_protection {
255 ipc_handler.enable_pii_protection()?;
256 }
257 let mcp_client = self.mcp_client.clone();
258 let builtin_executor = self.builtin_executor.clone();
259 let builtin_names = match &builtin_executor {
260 Some(exec) => exec
261 .list_builtin_tools()
262 .ok()
263 .map(|tools| tools.iter().map(|t| t.name.clone()).collect::<std::collections::HashSet<_>>()),
264 None => None,
265 };
266 let execution_timeout = Duration::from_secs(self.config.timeout_secs);
267
268 let ipc_task: JoinHandle<Result<()>> = tokio::spawn(async move {
269 let ipc_start = Instant::now();
270
271 while let Some(remaining_timeout) = execution_timeout.checked_sub(ipc_start.elapsed()) {
272 let Some(mut request) = ipc_handler.wait_for_request(remaining_timeout).await? else {
273 break;
274 };
275
276 debug!(
277 tool_name = %request.tool_name,
278 request_id = %request.id,
279 "Processing tool request from code"
280 );
281
282 if let Err(e) = ipc_handler.process_request_for_pii(&mut request) {
284 debug!(error = %e, "PII tokenization failed");
285 let response = ToolResponse {
286 id: request.id,
287 success: false,
288 result: None,
289 error: Some(format!("PII processing error: {e}")),
290 duration_ms: None,
291 cache_hit: None,
292 };
293 ipc_handler.write_response(response).await?;
294 continue;
295 }
296
297 let is_builtin = builtin_names
300 .as_ref()
301 .is_some_and(|set| set.contains(request.tool_name.as_str()));
302 let result = if is_builtin {
303 match builtin_executor
304 .as_ref()
305 .expect("builtin_names implies a builtin executor")
306 .execute_builtin_tool(&request.tool_name, &request.args)
307 .await
308 {
309 Ok(result) => {
310 debug!(tool_name = %request.tool_name, "Built-in tool executed successfully");
311 ToolResponse {
312 id: request.id,
313 success: true,
314 result: Some(result),
315 error: None,
316 duration_ms: None,
317 cache_hit: None,
318 }
319 }
320 Err(e) => {
321 debug!(tool_name = %request.tool_name, error = %e, "Built-in tool execution failed");
322 ToolResponse {
323 id: request.id,
324 success: false,
325 result: None,
326 error: Some(e.to_string()),
327 duration_ms: None,
328 cache_hit: None,
329 }
330 }
331 }
332 } else {
333 match mcp_client.execute_mcp_tool(&request.tool_name, &request.args).await {
334 Ok(result) => {
335 debug!(tool_name = %request.tool_name, "Tool executed successfully");
336 ToolResponse {
337 id: request.id,
338 success: true,
339 result: Some(result),
340 error: None,
341 duration_ms: None,
342 cache_hit: None,
343 }
344 }
345 Err(e) => {
346 debug!(
347 tool_name = %request.tool_name,
348 error = %e,
349 "Tool execution failed"
350 );
351 ToolResponse {
352 id: request.id,
353 success: false,
354 result: None,
355 error: Some(e.to_string()),
356 duration_ms: None,
357 cache_hit: None,
358 }
359 }
360 }
361 };
362
363 ipc_handler.write_response(result).await?;
365 }
366
367 Ok(())
368 });
369
370 let (program, args) = match self.language {
372 Language::Python3 => {
373 let interpreter = Language::detect_python_interpreter(&self.workspace_root);
374 if interpreter == "uv" {
375 (
377 "uv".to_string(),
378 vec![
379 "run".to_string(),
380 "python".to_string(),
381 code_file.to_string_lossy().into_owned(),
382 ],
383 )
384 } else {
385 (interpreter, vec![code_file.to_string_lossy().into_owned()])
386 }
387 }
388 Language::JavaScript => {
389 (self.language.interpreter().to_string(), vec![code_file.to_string_lossy().into_owned()])
390 }
391 };
392
393 let options = ProcessOptions {
394 program,
395 args,
396 env,
397 current_dir: Some(self.workspace_root.clone()),
398 timeout: Some(Duration::from_secs(self.config.timeout_secs)),
399 cancellation_token: None,
400 stdout: StreamCaptureConfig {
401 capture: true,
402 max_bytes: self.config.max_output_bytes,
403 },
404 stderr: StreamCaptureConfig {
405 capture: true,
406 max_bytes: self.config.max_output_bytes,
407 },
408 };
409
410 let process_output = AsyncProcessRunner::run(options).await.context("failed to execute code")?;
411
412 let duration_ms = start.elapsed().as_millis();
413
414 let stdout = String::from_utf8_lossy(&process_output.stdout).into_owned();
416 let stderr = String::from_utf8_lossy(&process_output.stderr).into_owned();
417
418 let json_result = self.extract_json_result(&stdout, self.language)?;
420
421 let _ = tokio::fs::remove_file(&code_file).await;
423 let _ = tokio::fs::remove_dir_all(&ipc_dir).await;
424
425 let ipc_result = async_utils::with_timeout(ipc_task, Duration::from_secs(1), "IPC handler task").await;
427
428 if let Err(e) = ipc_result {
429 debug!(error = %e, "IPC handler did not complete in time");
430 }
431
432 debug!(
433 exit_code = process_output.exit_status.code().unwrap_or(-1),
434 duration_ms,
435 has_json_result = json_result.is_some(),
436 "Code execution completed"
437 );
438
439 Ok(ExecutionResult {
440 exit_code: process_output.exit_status.code().unwrap_or(-1),
441 stdout,
442 stderr,
443 json_result,
444 duration_ms,
445 })
446 }
447
448 fn prepare_python_code(&self, sdk: &str, user_code: &str) -> Result<String> {
450 Ok(format!(
451 "{sdk}\n\n# User code\n{user_code}\n\n# Capture result\nimport json\nif 'result' in dir():\n print('__JSON_RESULT__')\n print(json.dumps(result, default=str))\n print('__END_JSON__')"
452 ))
453 }
454
455 fn prepare_javascript_code(&self, sdk: &str, user_code: &str) -> Result<String> {
457 Ok(format!(
458 "{sdk}\n\n// User code\n(async () => {{\n{user_code}\n\n// Capture result\nif (typeof result !== 'undefined') {{\n console.log('__JSON_RESULT__');\n console.log(JSON.stringify(result, null, 2));\n console.log('__END_JSON__');\n}}\n}})();\n"
459 ))
460 }
461
462 fn extract_json_result(&self, stdout: &str, _language: Language) -> Result<Option<Value>> {
464 if !stdout.contains("__JSON_RESULT__") {
465 return Ok(None);
466 }
467
468 let start_marker = "__JSON_RESULT__";
469 let end_marker = "__END_JSON__";
470
471 let start = match stdout.find(start_marker) {
472 Some(pos) => pos + start_marker.len(),
473 None => return Ok(None),
474 };
475
476 let end = match stdout[start..].find(end_marker) {
477 Some(pos) => start + pos,
478 None => return Ok(None),
479 };
480
481 let json_str = stdout[start..end].trim();
482
483 match serde_json::from_str::<Value>(json_str) {
484 Ok(value) => {
485 debug!("Extracted JSON result from code output");
486 Ok(Some(value))
487 }
488 Err(e) => {
489 debug!(error = %e, "Failed to parse JSON result");
490 Ok(None)
491 }
492 }
493 }
494
495 pub async fn generate_sdk(&self) -> Result<String> {
497 match self.language {
498 Language::Python3 => self.generate_python_sdk().await,
499 Language::JavaScript => self.generate_javascript_sdk().await,
500 }
501 }
502
503 async fn generate_python_sdk(&self) -> Result<String> {
505 debug!("Generating Python SDK for MCP tools");
506
507 let tools = self.mcp_client.list_mcp_tools().await.context("failed to list MCP tools")?;
508
509 let mut sdk = String::from(
510 r#"# MCP Tools SDK - Auto-generated
511import json
512import sys
513import os
514import time
515from typing import Any, Dict, Optional
516from uuid import uuid4
517
518class MCPTools:
519 """Interface to MCP tools from agent code via file-based IPC."""
520
521 IPC_DIR = os.environ.get("VTCODE_IPC_DIR")
522
523 def __init__(self):
524 self._call_count = 0
525 self._results = []
526 if not self.IPC_DIR:
527 raise RuntimeError("VTCODE_IPC_DIR is required for the VT Code SDK")
528 os.makedirs(self.IPC_DIR, exist_ok=True)
529
530 def _call_tool(self, name: str, args: Dict[str, Any]) -> Any:
531 """Call an MCP tool via file-based IPC."""
532 request_id = str(uuid4())
533
534 # Write request
535 request = {
536 "id": request_id,
537 "tool_name": name,
538 "args": args
539 }
540 request_file = os.path.join(self.IPC_DIR, "request.json")
541 with open(request_file, 'w') as f:
542 json.dump(request, f)
543
544 # Wait for response
545 response_file = os.path.join(self.IPC_DIR, "response.json")
546 timeout = 30
547 start = time.time()
548 while time.time() - start < timeout:
549 if os.path.exists(response_file):
550 with open(response_file, 'r') as f:
551 response = json.load(f)
552
553 if response.get("id") == request_id:
554 # Clean up response
555 try:
556 os.remove(response_file)
557 except:
558 pass
559
560 if response.get("success"):
561 return response.get("result")
562 else:
563 raise RuntimeError(f"Tool error: {response.get('error', 'unknown error')}")
564
565 time.sleep(0.1)
566
567 raise TimeoutError(f"Tool '{name}' timed out after {timeout}s")
568
569 def log(self, message: str) -> None:
570 """Log a message that will be captured."""
571 print(f"[LOG] {message}")
572
573# Initialize tools interface
574mcp = MCPTools()
575"#,
576 );
577
578 for tool in tools {
580 sdk.push_str(&format!(
581 "\ndef {}(**kwargs):\n \"\"\"{}.\"\"\"\n return mcp._call_tool('{}', kwargs)\n\n",
582 sanitize_function_name(&tool.name),
583 tool.description,
584 tool.name
585 ));
586 }
587
588 if let Some(builtin) = &self.builtin_executor {
590 if let Ok(builtin_tools) = builtin.list_builtin_tools() {
591 for tool in builtin_tools {
592 sdk.push_str(&format!(
593 "\ndef {}(**kwargs):\n \"\"\"{}.\"\"\"\n return mcp._call_tool('{}', kwargs)\n\n",
594 sanitize_function_name(&tool.name),
595 tool.description,
596 tool.name
597 ));
598 }
599 }
600 }
601
602 Ok(sdk)
603 }
604
605 async fn generate_javascript_sdk(&self) -> Result<String> {
607 debug!("Generating JavaScript SDK for MCP tools");
608
609 let tools = self.mcp_client.list_mcp_tools().await.context("failed to list MCP tools")?;
610
611 let mut sdk = String::from(
612 r#"// MCP Tools SDK - Auto-generated
613const fs = require('fs');
614const path = require('path');
615const { v4: uuid4 } = require('uuid');
616
617class MCPTools {
618 constructor() {
619 this.callCount = 0;
620 this.results = [];
621 this.ipcDir = process.env.VTCODE_IPC_DIR;
622 if (!this.ipcDir) {
623 throw new Error('VTCODE_IPC_DIR is required for the VT Code SDK');
624 }
625 if (!fs.existsSync(this.ipcDir)) {
626 fs.mkdirSync(this.ipcDir, { recursive: true });
627 }
628 }
629
630 async callTool(name, args = {}) {
631 const requestId = uuid4();
632 const request = {
633 id: requestId,
634 tool_name: name,
635 args: args
636 };
637
638 const requestFile = path.join(this.ipcDir, 'request.json');
639 fs.writeFileSync(requestFile, JSON.stringify(request, null, 2));
640
641 // Wait for response
642 const responseFile = path.join(this.ipcDir, 'response.json');
643 const timeout = 30000; // 30s
644 const start = Date.now();
645
646 while (Date.now() - start < timeout) {
647 try {
648 if (fs.existsSync(responseFile)) {
649 const response = JSON.parse(fs.readFileSync(responseFile, 'utf-8'));
650
651 if (response.id === requestId) {
652 // Clean up response
653 try {
654 fs.unlinkSync(responseFile);
655 } catch (e) {}
656
657 if (response.success) {
658 return response.result;
659 } else {
660 throw new Error(`Tool error: ${response.error || 'unknown error'}`);
661 }
662 }
663 }
664 } catch (e) {
665 if (e.code !== 'ENOENT') throw e;
666 }
667
668 await new Promise(r => setTimeout(r, 100));
669 }
670
671 throw new Error(`Tool '${name}' timed out after ${timeout}ms`);
672 }
673
674 log(message) {
675 console.log(`[LOG] ${message}`);
676 }
677}
678
679const mcp = new MCPTools();
680
681"#,
682 );
683
684 for tool in tools {
686 sdk.push_str(&format!(
687 "async function {}(args = {{}}) {{\n // {}\n return await mcp.callTool('{}', args);\n}}\n\n",
688 sanitize_function_name(&tool.name),
689 tool.description,
690 tool.name
691 ));
692 }
693
694 if let Some(builtin) = &self.builtin_executor {
696 if let Ok(builtin_tools) = builtin.list_builtin_tools() {
697 for tool in builtin_tools {
698 sdk.push_str(&format!(
699 "async function {}(args = {{}}) {{\\n // {}\\n return await mcp.callTool('{}', args);\\n}}\\n\\n",
700 sanitize_function_name(&tool.name), tool.description, tool.name
701 ));
702 }
703 }
704 }
705
706 Ok(sdk)
707 }
708}
709
710fn sanitize_function_name(name: &str) -> String {
712 name.chars()
713 .map(|c| if c.is_ascii_alphanumeric() || c == '_' { c } else { '_' })
714 .collect()
715}
716
717#[cfg(test)]
718mod tests {
719 use super::*;
720 use crate::mcp::{McpClientStatus, McpToolExecutor, McpToolInfo};
721 use async_trait::async_trait;
722 use serde_json::{Value, json};
723 use std::path::PathBuf;
724 use std::sync::Arc;
725
726 struct MockMcpToolExecutor;
727
728 #[async_trait]
729 impl McpToolExecutor for MockMcpToolExecutor {
730 async fn execute_mcp_tool(&self, _tool_name: &str, _args: &Value) -> Result<Value> {
731 Ok(json!({}))
732 }
733
734 async fn list_mcp_tools(&self) -> Result<Vec<McpToolInfo>> {
735 Ok(vec![McpToolInfo {
736 name: "read_file".to_string(),
737 description: "Read a file".to_string(),
738 provider: "test".to_string(),
739 input_schema: json!({}),
740 output_schema: None,
741 }])
742 }
743
744 async fn has_mcp_tool(&self, _tool_name: &str) -> Result<bool> {
745 Ok(true)
746 }
747
748 fn get_status(&self) -> McpClientStatus {
749 McpClientStatus {
750 enabled: true,
751 provider_count: 1,
752 active_connections: 1,
753 configured_providers: vec!["test".to_string()],
754 }
755 }
756 }
757
758 fn test_executor(language: Language) -> CodeExecutor {
759 CodeExecutor::new(language, Arc::new(MockMcpToolExecutor), PathBuf::from("/workspace"))
760 }
761
762 struct MockBuiltinExecutor;
763
764 #[async_trait]
765 impl BuiltinToolExecutor for MockBuiltinExecutor {
766 async fn execute_builtin_tool(&self, _name: &str, _args: &Value) -> Result<Value> {
767 Ok(json!({}))
768 }
769
770 fn list_builtin_tools(&self) -> Result<Vec<BuiltinToolInfo>> {
771 Ok(vec![BuiltinToolInfo {
772 name: "unified_file".to_string(),
773 description: "Read/write files".to_string(),
774 }])
775 }
776 }
777
778 #[tokio::test]
779 async fn generate_sdk_includes_builtin_wrappers() {
780 let executor = test_executor(Language::Python3).with_builtin_executor(Some(Arc::new(MockBuiltinExecutor)));
781 let sdk = executor.generate_python_sdk().await.unwrap();
782 assert!(sdk.contains("def unified_file("));
783 assert!(sdk.contains("mcp._call_tool('unified_file'"));
784 assert!(sdk.contains("def read_file("));
786 }
787
788 #[test]
789 fn sanitize_function_name_handles_special_chars() {
790 assert_eq!(sanitize_function_name("read_file"), "read_file");
791 assert_eq!(sanitize_function_name("read-file"), "read_file");
792 assert_eq!(sanitize_function_name("read.file"), "read_file");
793 assert_eq!(sanitize_function_name("readFile123"), "readFile123");
794 }
795
796 #[test]
797 fn language_as_str() {
798 assert_eq!(Language::Python3.as_str(), "python3");
799 assert_eq!(Language::JavaScript.as_str(), "javascript");
800 }
801
802 #[test]
803 fn language_interpreter() {
804 assert_eq!(Language::Python3.interpreter(), "python3");
805 assert_eq!(Language::JavaScript.interpreter(), "node");
806 }
807
808 #[tokio::test]
809 async fn python_sdk_uses_ipc_dir_only() {
810 let sdk = test_executor(Language::Python3)
811 .generate_sdk()
812 .await
813 .expect("python sdk should generate");
814
815 assert!(sdk.contains("VTCODE_IPC_DIR"));
816 assert!(!sdk.contains("/tmp/vtcode_ipc"));
817 assert!(!sdk.contains("VTCODE_WORKSPACE"));
818 }
819
820 #[tokio::test]
821 async fn javascript_sdk_uses_ipc_dir_only() {
822 let sdk = test_executor(Language::JavaScript)
823 .generate_sdk()
824 .await
825 .expect("javascript sdk should generate");
826
827 assert!(sdk.contains("VTCODE_IPC_DIR"));
828 assert!(!sdk.contains("/tmp/vtcode_ipc"));
829 assert!(!sdk.contains("VTCODE_WORKSPACE"));
830 }
831}