| """Visual browser agent - Chrome with live stream and agent memory. |
|
|
| Uses E2B Desktop sandbox with Chrome browser.
|
| Time limit: 5 minutes (300 seconds)
|
| Shows live video stream.
|
| Includes full memory/history tracking via AgentState.
|
| """
|
|
|
| import os |
| import shlex |
| import logging |
| import time |
| import json |
| from typing import AsyncGenerator, Optional |
|
|
| from app.agents.browser_dom import build_visual_dom_extract_script |
| from app.agents.browser_search import choose_search_url |
| from app.agents.browser_decision import decide_browser_action |
| from app.config import get_settings |
| from app.agents.llm_client import generate_completion |
| from app.agents.graph.state import AgentState |
| from app.agents.flaresolverr import is_cloudflare_blocked |
|
|
| logger = logging.getLogger(__name__)
|
|
|
| MAX_TIME_SECONDS = 300
|
|
|
|
|
| async def run_browser_visual_agent( |
| task: str, |
| url: Optional[str] = None, |
| ) -> AsyncGenerator[dict, None]: |
| """Run the visual browser agent with Chrome and live stream.""" |
| settings = get_settings()
|
|
|
| if not settings.e2b_api_key:
|
| yield {"type": "error", "message": "E2B_API_KEY not configured"}
|
| return
|
|
|
|
|
| state = AgentState(
|
| task=task,
|
| url=url,
|
| timeout_seconds=MAX_TIME_SECONDS,
|
| start_time=time.time()
|
| )
|
|
|
| yield {"type": "status", "message": "🚀 Initializing agent..."}
|
|
|
| desktop = None
|
|
|
| try:
|
| from e2b_desktop import Sandbox
|
|
|
| os.environ["E2B_API_KEY"] = settings.e2b_api_key
|
|
|
| yield {"type": "status", "message": "🖥️ Creating virtual desktop..."}
|
| desktop = Sandbox.create(timeout=600)
|
| state.desktop = desktop
|
|
|
|
|
| stream_url = None
|
| try:
|
| desktop.stream.start(require_auth=True)
|
| auth_key = desktop.stream.get_auth_key()
|
| stream_url = desktop.stream.get_url(auth_key=auth_key)
|
| yield {"type": "stream", "url": stream_url}
|
| logger.info(f"Stream started: {stream_url}")
|
| desktop.wait(2000)
|
| except Exception as e:
|
| logger.warning(f"Could not start stream: {e}")
|
|
|
|
|
| yield {"type": "status", "message": "🌐 Launching browser..."}
|
|
|
| if url: |
| start_url = url |
| else: |
| start_url = choose_search_url(task, visited_urls=state.visited_urls) |
| state.add_query(task) |
|
|
| chrome_flags = "--no-sandbox --disable-gpu --start-maximized --no-first-run --disable-default-apps --disable-popup-blocking --disable-translate --no-default-browser-check"
|
| desktop.commands.run(f"google-chrome {chrome_flags} {shlex.quote(start_url)} &", background=True)
|
| desktop.wait(3000)
|
|
|
|
|
| desktop.press("enter")
|
| desktop.wait(1000)
|
|
|
|
|
| state.visited_urls.append(start_url)
|
| state.add_action({"type": "navigate", "url": start_url})
|
|
|
|
|
| while state.should_continue(): |
| state.step_count += 1 |
| elapsed = int(state.get_elapsed_time()) |
| remaining = int(state.get_remaining_time()) |
|
|
| yield {"type": "status", "message": f"🔍 Step {state.step_count}: Analyzing... ({elapsed}s / {MAX_TIME_SECONDS}s)"}
|
|
|
| |
| current_url = state.visited_urls[-1] |
| page_content = "" |
| page_links: list[str] = [] |
| extracted_blocked = False |
| |
| try: |
| script = build_visual_dom_extract_script(current_url) |
| desktop.commands.run(f"cat > /tmp/visual_dom_extract.py << 'EOF'\n{script}\nEOF", timeout=10) |
| result = desktop.commands.run("python3 /tmp/visual_dom_extract.py", timeout=45) |
| output = result.stdout.strip() if hasattr(result, "stdout") else "" |
| data = json.loads(output) if output else {} |
| page_content = str(data.get("content", "") or "") |
| page_links = [ |
| link for link in (data.get("links", []) or []) |
| if isinstance(link, str) and link.startswith("http") |
| ] |
| extracted_blocked = bool(data.get("blocked", False)) |
| state.page_content = page_content |
| if data.get("error"): |
| state.add_error(f"DOM extraction warning: {data['error']}") |
| except Exception as e: |
| logger.warning(f"DOM extraction failed: {e}") |
| state.add_error(f"DOM extraction failed: {e}") |
|
|
| preview_text = page_content[:2000] if page_content else "(empty page)" |
| |
| |
| is_blocked = extracted_blocked or (is_cloudflare_blocked(page_content) if page_content else False) |
| |
| if is_blocked: |
| yield {"type": "status", "message": f"🚫 Cloudflare at {current_url[:40]}..., trying next link..."}
|
| state.add_error(f"Cloudflare blocked: {current_url}")
|
| else:
|
|
|
| state.extracted_data.append({ |
| "url": current_url, |
| "content_length": len(page_content), |
| "links_found": len(page_links), |
| "preview": page_content[:200] |
| }) |
| |
| decision = await decide_browser_action( |
| task=task, |
| current_url=current_url, |
| state=state, |
| content_preview=preview_text, |
| blocked=is_blocked, |
| allow_scroll=False, |
| mode_label="visual Chrome", |
| step_label=f"{state.step_count}, {remaining}s remaining", |
| links=page_links, |
| max_tokens=600, |
| ) |
| |
| action = decision.get("action", "DONE") |
| value = decision.get("value", "") |
| final_answer = decision.get("answer", "") |
| reason = decision.get("reason", "") |
| known_facts = decision.get("known_facts", []) |
| missing_points = decision.get("missing_points", []) |
|
|
| if action == "SEARCH": |
| state.add_query(value) |
|
|
| if isinstance(known_facts, list) or isinstance(missing_points, list): |
| state.update_research_progress( |
| known_facts=known_facts if isinstance(known_facts, list) else None, |
| missing_points=missing_points if isinstance(missing_points, list) else None, |
| ) |
|
|
|
|
| state.add_action({"type": action.lower(), "value": value, "reason": reason})
|
|
|
| yield {"type": "status", "message": f"🤔 Action: {action} - {reason[:50]}"}
|
|
|
| yield { |
| "type": "progress", |
| "known_facts": state.known_facts[-8:], |
| "missing_points": state.missing_points[-8:], |
| "last_queries": state.last_queries[-8:], |
| } |
|
|
| if action == "DONE": |
| state.success = True |
| |
| if not final_answer: |
| |
| all_content = "\n\n".join([ |
| f"Source: {d['url']}\n{d.get('preview', '')}" |
| for d in state.extracted_data[-5:] |
| ]) |
| known_summary = "\n".join([f"- {f}" for f in state.known_facts[-8:]]) or "(none)" |
| missing_summary = "\n".join([f"- {m}" for m in state.missing_points[-8:]]) or "(none)" |
| final_prompt = ( |
| f"Based on this content, answer: {task}\n\n" |
| f"Known facts:\n{known_summary}\n\n" |
| f"Missing points:\n{missing_summary}\n\n" |
| f"Content:\n{all_content}" |
| ) |
| final_answer = await generate_completion( |
| messages=[{"role": "user", "content": final_prompt}], |
| max_tokens=1000 |
| ) |
|
|
| state.final_result = final_answer
|
|
|
| yield {"type": "stream_end", "message": "Done"}
|
| yield {
|
| "type": "result",
|
| "content": final_answer,
|
| "links": state.visited_urls,
|
| "steps": state.step_count,
|
| "success": True
|
| }
|
|
|
| yield {"type": "complete", "message": f"Completed in {int(state.get_elapsed_time())}s with {state.step_count} steps"}
|
| return
|
|
|
| elif action == "SEARCH": |
| new_url = choose_search_url( |
| value, |
| visited_urls=state.visited_urls, |
| current_url=current_url, |
| blocked=is_blocked, |
| ) |
| |
| if new_url not in state.visited_urls: |
| desktop.commands.run(f"google-chrome {shlex.quote(new_url)} &", background=True) |
| desktop.wait(3000) |
| state.visited_urls.append(new_url) |
|
|
| elif action == "NAVIGATE":
|
| if value and value.startswith("http"):
|
| if value in state.visited_urls:
|
| yield {"type": "status", "message": f"⏭️ Already visited, skipping..."}
|
| state.add_error(f"Tried to revisit: {value}")
|
| else:
|
| desktop.commands.run(f"google-chrome {shlex.quote(value)} &", background=True)
|
| desktop.wait(3000)
|
| state.visited_urls.append(value)
|
|
|
| |
| desktop.wait(1000) |
|
|
|
|
| yield {"type": "status", "message": "⏰ Time limit reached, generating final answer from memory..."}
|
|
|
| all_content = "\n\n".join([ |
| f"Source: {d['url']}\n{d.get('preview', '')}" |
| for d in state.extracted_data[-5:] |
| ]) |
| known_summary = "\n".join([f"- {f}" for f in state.known_facts[-8:]]) or "(none)" |
| missing_summary = "\n".join([f"- {m}" for m in state.missing_points[-8:]]) or "(none)" |
| final_prompt = ( |
| f"Based on this content, answer: {task}\n\n" |
| f"Known facts:\n{known_summary}\n\n" |
| f"Missing points:\n{missing_summary}\n\n" |
| f"Content:\n{all_content}" |
| ) |
| final_answer = await generate_completion( |
| messages=[{"role": "user", "content": final_prompt}], |
| max_tokens=1000 |
| ) |
|
|
| state.final_result = final_answer
|
|
|
| yield {"type": "stream_end", "message": "Done"}
|
| yield {
|
| "type": "result",
|
| "content": final_answer,
|
| "links": state.visited_urls,
|
| "steps": state.step_count,
|
| "success": True
|
| }
|
| yield {"type": "complete", "message": f"Completed in {MAX_TIME_SECONDS}s (timeout) with {state.step_count} steps"}
|
|
|
| except ImportError as e:
|
| yield {"type": "error", "message": "e2b-desktop not installed"}
|
| except Exception as e:
|
| logger.exception("Browser agent error")
|
| yield {"type": "error", "message": f"Error: {str(e)}"}
|
| finally:
|
| if desktop:
|
| try:
|
| desktop.stream.stop()
|
| except:
|
| pass
|
| try:
|
| desktop.kill()
|
| except:
|
| pass
|
|
|