canopy-258m-r3 / miniswardbower /browser /swarm_coordinator.py
psikosen's picture
Update to v6: Tri-Engine Swarm Coordinator, URL Invariance, Thought Bus & State Verification
ed79a7f verified
Raw History Blame Contribute Delete
7.82 kB
"""
Cooperative Tri-Engine Browser Swarm Coordinator for Miniswardbower v6.
Orchestrates specialized browser execution engines tailored to task modalities:
1. Lightpanda: Ultra-fast headless Zig engine (~27.5MB RSS, ~15ms load) for
rapid link scraping, structural AXTree indexing, and text distillation.
2. Obscura: Native Rust stealth CDP engine (~49.3MB RSS, authentic fingerprinting)
with native rasterization for Set-of-Marks visual auditing and form fills.
3. Chromium: Full Blink fidelity fallback for complex client-side SPAs.
"""
from __future__ import annotations
import asyncio
import logging
import time
from typing import Any, Dict, List, Optional, Tuple, Union
from miniswardbower.browser.controller import BrowserController
from miniswardbower.core.config import BrowserConfig
from miniswardbower.core.schemas import (
ActionBatch,
BrowserAction,
BrowserActionType,
PrunedAXTree,
)
logger = logging.getLogger(__name__)
class TriEngineSwarmCoordinator:
"""
Cooperative multi-engine browser swarm manager.
Dynamically routes browser interactions to the optimal specialized engine:
- Indexing / Scraping -> Lightpanda (maximum speed, zero rasterization waste)
- Visual Audits / Stealth Actions -> Obscura (anti-bot fingerprinting + native screenshots)
- Deep SPA / Fallback -> Chromium (complete desktop Blink compatibility)
"""
def __init__(self, base_config: Optional[BrowserConfig] = None):
self.base_config = base_config or BrowserConfig()
self._controllers: Dict[str, BrowserController] = {}
self._active_engine: str = self.base_config.browser_engine
self._telemetry: Dict[str, Any] = {
"lightpanda": {"actions": 0, "total_time_ms": 0.0, "rss_mb": 0.0},
"obscura": {"actions": 0, "total_time_ms": 0.0, "rss_mb": 0.0},
"chromium": {"actions": 0, "total_time_ms": 0.0, "rss_mb": 0.0},
}
async def get_controller(self, engine: str) -> BrowserController:
"""Lazily launches and returns the controller for the specified engine."""
if engine not in self._controllers:
cfg = self.base_config.model_copy()
cfg.browser_engine = engine
ctrl = BrowserController(cfg)
await ctrl.start()
self._controllers[engine] = ctrl
logger.info("Initialized swarm controller for engine '%s'", engine)
return self._controllers[engine]
async def scrape_and_index(
self, url: str
) -> Tuple[PrunedAXTree, str]:
"""
High-throughput indexing pass powered by Lightpanda.
Extracts structural accessibility tree and visible text in ~15ms.
"""
t0 = time.perf_counter()
engine = "lightpanda"
try:
ctrl = await self.get_controller(engine)
except Exception as e:
logger.warning("Lightpanda unavailable (%s); falling back to Obscura", e)
engine = "obscura"
ctrl = await self.get_controller(engine)
await ctrl.goto(url)
tree = await ctrl.get_pruned_tree()
scraped = await ctrl.scrape_page()
text = scraped.get("body_text", "") if isinstance(scraped, dict) else ""
if not text:
try:
text = await ctrl.page.inner_text("body")
except Exception:
text = ""
elapsed = (time.perf_counter() - t0) * 1000.0
self._telemetry[engine]["actions"] += 1
self._telemetry[engine]["total_time_ms"] += elapsed
self._active_engine = engine
return tree, text
async def visual_interact_and_submit(
self,
url: str,
actions: List[BrowserAction],
capture_audit_screenshot: bool = True,
) -> Tuple[List[Dict[str, Any]], Optional[bytes]]:
"""
Stealth interaction pass powered by Obscura.
Executes form inputs, clicks, and captures visual Set-of-Marks audit screenshots.
"""
t0 = time.perf_counter()
engine = "obscura"
try:
ctrl = await self.get_controller(engine)
except Exception as e:
logger.warning("Obscura unavailable (%s); falling back to Chromium", e)
engine = "chromium"
ctrl = await self.get_controller(engine)
# Navigate to target page if different from current
if ctrl.page.url != url:
await ctrl.goto(url)
# Execute actions
tree = await ctrl.get_pruned_tree()
results = await ctrl.execute_chunk(actions, tree=tree)
screenshot_bytes = None
if capture_audit_screenshot:
try:
screenshot_bytes = await ctrl.page.screenshot()
except Exception as sc_err:
logger.warning("Audit screenshot capture failed: %s", sc_err)
elapsed = (time.perf_counter() - t0) * 1000.0
self._telemetry[engine]["actions"] += len(actions)
self._telemetry[engine]["total_time_ms"] += elapsed
self._active_engine = engine
return results, screenshot_bytes
async def execute_cooperative_pipeline(
self,
start_url: str,
form_actions: List[BrowserAction],
) -> Dict[str, Any]:
"""
Executes a 2-stage cooperative pipeline:
Stage 1: Lightpanda rapidly maps the page structure and locates element selectors.
Stage 2: Obscura executes actions with authentic anti-bot fingerprints and visual receipts.
"""
t_start = time.perf_counter()
# Stage 1: Fast Scraping
tree, text = await self.scrape_and_index(start_url)
t_scrape = (time.perf_counter() - t_start) * 1000.0
# Stage 2: Interaction & Receipt
t_act_start = time.perf_counter()
results, screenshot = await self.visual_interact_and_submit(
start_url, form_actions, capture_audit_screenshot=True
)
t_interact = (time.perf_counter() - t_act_start) * 1000.0
t_total = (time.perf_counter() - t_start) * 1000.0
return {
"elements_indexed": len(tree.elements),
"text_length": len(text),
"actions_executed": len(results),
"action_results": results,
"has_visual_receipt": screenshot is not None,
"receipt_bytes": len(screenshot) if screenshot else 0,
"timings_ms": {
"scrape_stage_ms": round(t_scrape, 2),
"interact_stage_ms": round(t_interact, 2),
"total_ms": round(t_total, 2),
},
"telemetry": self.get_telemetry(),
}
def get_telemetry(self) -> Dict[str, Any]:
"""Returns runtime stats and memory usage across all initialized engines."""
report = {}
for engine, ctrl in self._controllers.items():
rss = 0.0
driver = getattr(ctrl, f"_{engine}_driver", None)
if driver:
mgr = getattr(driver, "server_manager", None) or getattr(driver, "server_mgr", None)
if mgr and hasattr(mgr, "get_memory_rss_mb"):
rss = mgr.get_memory_rss_mb()
report[engine] = {
"actions": self._telemetry[engine]["actions"],
"total_time_ms": round(self._telemetry[engine]["total_time_ms"], 2),
"rss_mb": round(rss, 2),
}
return report
async def close_all(self) -> None:
"""Gracefully shuts down all managed engine controllers and background daemons."""
for engine, ctrl in list(self._controllers.items()):
try:
await ctrl.stop()
except Exception as e:
logger.error("Error stopping engine %s: %s", engine, e)
self._controllers.clear()