Spaces:
Running
Running
File size: 8,973 Bytes
09801ca | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 | import json
import logging
import os
import re
from typing import Dict, Any, List, AsyncGenerator
from core.llm import chat
logger = logging.getLogger(__name__)
class AgentEngine:
"""
An XML-based ReAct agent engine that streams Server-Sent Events (SSE) to the frontend.
Supports tools: read_file, write_file, run_command, ask_user.
"""
def __init__(self, workspace_dir: str, user_id: str):
self.workspace_dir = workspace_dir
self.user_id = user_id
async def stream_chat(self, prompt: str, history: List[Dict[str, str]], files: Dict[str, str], model: str) -> AsyncGenerator[str, None]:
# Pre-process file summary
file_summary = ""
for name, content in files.items():
file_summary += f"- {name} ({len(content)} bytes)\n"
system_prompt = f"""You are DataVision Agent, a highly capable Silicon Valley AI Architect.
You operate in a continuous Loop: THOUGHT -> ACTION -> OBSERVATION.
CRITICAL RULE: You are fully capable of writing files and executing commands on the user's local system using the tools below. NEVER say you cannot write files or ask the user to implement things manually. You MUST use your tools to complete the user's task yourself!
CRITICAL RULE 2: If a file or project does not exist, DO NOT complain to the user! You must CREATE the files yourself using `write_file` or initialize the project yourself using `run_command` (e.g. `npm init`, `npx create-vite`, etc.). Be proactive!
Available Tools:
1. `read_file`: Read the contents of a file. Args: {{"path": "string", "offset": "optional line start (int)", "limit": "optional max lines (int)"}}
2. `write_file`: Write or overwrite a file.
3. `run_command`: Run a terminal command.
4. `ask_user`: Ask the user a clarifying question or request permission to proceed.
5. `mcp_call`: Call a Model Context Protocol (MCP) server tool. Args: {{"server": "name", "tool": "tool_name", "args": {{}}}}
To use a tool, you MUST use XML tags exactly like this:
<thought>
I need to check the package.json to see if react is installed.
</thought>
<tool name="read_file">
{{"path": "package.json"}}
</tool>
When you are finished and want to respond to the user, use:
<thought>
I have completed the task.
</thought>
<response>
Your task is complete! Here is what I did...
</response>
Workspace Context:
{file_summary}
"""
# We will loop up to 10 times autonomously
MAX_TURNS = 10
current_messages = list(history) if history else []
current_messages.append({"role": "user", "content": prompt})
for turn in range(MAX_TURNS):
yield f"data: {json.dumps({'type': 'status', 'content': f'Turn {turn+1}/{MAX_TURNS}: Thinking...'})}\n\n"
try:
# Call LLM
response_text = chat(
messages=current_messages,
system=system_prompt,
model=model,
temperature=0.2
)
# We append assistant's raw output to history
current_messages.append({"role": "assistant", "content": response_text})
# Parse the response for thoughts, tools, and final response
thought_match = re.search(r'<thought>(.*?)</thought>', response_text, re.DOTALL)
if thought_match:
thought = thought_match.group(1).strip()
yield f"data: {json.dumps({'type': 'thought', 'content': thought})}\n\n"
tool_match = re.search(r'<tool name="(.*?)">(.*?)</tool>', response_text, re.DOTALL)
response_match = re.search(r'<response>(.*?)</response>', response_text, re.DOTALL)
if response_match:
# The agent is done
final_text = response_match.group(1).strip()
yield f"data: {json.dumps({'type': 'message', 'content': final_text})}\n\n"
break
elif tool_match:
tool_name = tool_match.group(1).strip()
try:
tool_args = json.loads(tool_match.group(2).strip())
if not isinstance(tool_args, dict):
tool_args = {"raw": str(tool_args)}
except:
tool_args = {"raw": tool_match.group(2).strip()}
if tool_name == "ask_user":
# Stream the question to the user and stop the loop
question = tool_args.get("question", tool_args.get("raw", "I need your input."))
yield f"data: {json.dumps({'type': 'message', 'content': question})}\n\n"
break
yield f"data: {json.dumps({'type': 'tool_call', 'tool': tool_name, 'args': tool_args})}\n\n"
# Execute Tool
observation = self._execute_tool(tool_name, tool_args)
yield f"data: {json.dumps({'type': 'tool_result', 'tool': tool_name, 'result': observation})}\n\n"
# Append observation and loop again
current_messages.append({"role": "user", "content": f"<observation>\n{observation}\n</observation>"})
else:
# Neither a tool nor a response was found. Assume they just responded directly.
yield f"data: {json.dumps({'type': 'message', 'content': response_text})}\n\n"
break
except Exception as e:
logger.error(f"Agent Loop Error: {e}")
yield f"data: {json.dumps({'type': 'error', 'content': str(e)})}\n\n"
break
# Sync files back to frontend
updated_files = {}
for root, _, files_in_dir in os.walk(self.workspace_dir):
for file in files_in_dir:
if file.endswith(('.py', '.ts', '.tsx', '.js', '.jsx', '.html', '.css', '.md', '.json', '.txt', '.yml', '.yaml', '.toml', '.sh', '.bat')):
full_path = os.path.join(root, file)
rel_path = os.path.relpath(full_path, self.workspace_dir)
rel_path = rel_path.replace("\\", "/") # Normalize for web
try:
with open(full_path, "r", encoding="utf-8") as f:
updated_files[rel_path] = f.read()
except:
pass
yield f"data: {json.dumps({'type': 'sync_files', 'files': updated_files})}\n\n"
yield f"data: {json.dumps({'type': 'done'})}\n\n"
def _execute_tool(self, name: str, args: dict) -> str:
try:
if name == "read_file":
path = os.path.join(self.workspace_dir, args.get("path", ""))
offset = args.get("offset")
limit = args.get("limit")
with open(path, "r", encoding="utf-8") as f:
lines = f.readlines()
if offset is not None:
lines = lines[int(offset):]
if limit is not None:
lines = lines[:int(limit)]
return "".join(lines)
elif name == "write_file":
path = os.path.join(self.workspace_dir, args.get("path", ""))
os.makedirs(os.path.dirname(path), exist_ok=True)
with open(path, "w", encoding="utf-8") as f:
f.write(args.get("content", ""))
return "File written successfully."
elif name == "run_command":
import subprocess
cmd = args.get("command", "")
env = os.environ.copy()
env["PYTHONIOENCODING"] = "utf-8"
result = subprocess.run(
cmd, shell=True, cwd=self.workspace_dir,
capture_output=True, text=True, encoding="utf-8", errors="replace", env=env, timeout=30
)
output = result.stdout + result.stderr
return output if output else "Command executed silently with code 0."
elif name == "mcp_call":
# Stub for MCP client.
# In a full implementation, this would connect to the MCP server and call the tool.
server = args.get("server")
tool = args.get("tool")
return f"MCP Error: Server '{server}' is not configured. Please configure MCP servers in settings."
else:
return f"Unknown tool: {name}"
except Exception as e:
return f"Error executing tool: {e}"
|