diff --git a/agent.py b/agent.py deleted file mode 100644 index bc70bd8..0000000 --- a/agent.py +++ /dev/null @@ -1,1001 +0,0 @@ -"""Private compliant website scraper using SeleniumBase/browser automation. - -The public skill is intentionally narrow: it accepts one starting URL, a CSS -selector mapping, pagination limits, browser options, robots.txt behavior, and -output format. It never attempts to bypass CAPTCHAs, anti-bot challenges, -access controls, paywalls, login gates, robots.txt restrictions, or terms of -service restrictions. - -This module now also exposes an LLM-powered iterative skill that can plan and -refine safe scraping runs using the base browser skill as an execution step. -It uses a platform-provided LLM grant (ctx.llm) with strict safety rails and -never attempts to evade protections. The deep agent proposes selector/URL -adjustments limited to the same origin and stops on any challenge. -""" -from __future__ import annotations - -import json -import re -import textwrap -import urllib.parse -import urllib.request -import urllib.robotparser -from datetime import datetime, timezone -from pathlib import Path -from typing import Any, Literal - -from pydantic import BaseModel, Field - -from a2a_pack import ( - A2AAgent, - EgressPolicy, - LLMProvisioning, - NoAuth, - Pricing, - Resources, - RunContext, - WorkspaceAccess, - WorkspaceMode, - skill, -) - - -OUTPUT_DIR = "outputs/seleniumbase-website-scraper" -RUNTIME_SKILLS_DIR = "seleniumbase-website-scraper/.deepagents/skills/" -DEEPAGENTS_RECURSION_LIMIT = 500 - - -class SeleniumbaseWebsiteScraperConfig(BaseModel): - max_pages_limit: int = Field(default=50, ge=1, le=200) - default_user_agent: str = "A2A seleniumbase-website-scraper/1.0 (+authorized public/user-owned scraping only)" - - -class SeleniumbaseWebsiteScraper(A2AAgent[SeleniumbaseWebsiteScraperConfig, NoAuth]): - name = "seleniumbase-website-scraper" - description = ( - "Private compliant website scraping agent using SeleniumBase/browser " - "automation with robots.txt checks, same-domain pagination safeguards, " - "screenshots, HTML capture, JSON/CSV output, retries, and polite rate limiting." - ) - version = "0.1.0" - - config_model = SeleniumbaseWebsiteScraperConfig - auth_model = NoAuth - - llm_provisioning = LLMProvisioning.PLATFORM - pricing = Pricing( - price_per_call_usd=0.0, - caller_pays_llm=True, - notes=( - "Browser scraping is deterministic; any LLM use reads the caller-provided ctx.llm credential. " - "The caller pays any forwarded LLM cost." - ), - ) - resources = Resources(cpu="2", memory="2Gi", max_runtime_seconds=900) - egress = EgressPolicy(deny_internet_by_default=False) - tools_used = ("seleniumbase", "selenium", "chromium", "browser-automation", "deepagents", "langchain") - workspace_access = WorkspaceAccess.dynamic( - max_files=256, - allowed_modes=(WorkspaceMode.READ_ONLY, WorkspaceMode.READ_WRITE_OVERLAY), - require_reason=False, - max_total_size_bytes=250 * 1024 * 1024, - ) - - @skill( - description=( - "Authorized website scraping with SeleniumBase/browser automation, " - "robots.txt checks, selector extraction, same-domain pagination, " - "screenshots, and JSON/CSV output. Stops on CAPTCHA, anti-bot, login, " - "paywall, or access-control challenges." - ), - timeout_seconds=900, - cost_class="browser-heavy", - grant_outputs_prefix=OUTPUT_DIR, - grant_write_prefixes=(OUTPUT_DIR,), - grant_run_timeout_seconds=900, - ) - async def scrape_website( - self, - ctx: RunContext[NoAuth], - url: str, - selectors: dict[str, str], - max_pages: int = 1, - wait_seconds: float = 3.0, - output_format: Literal["json", "csv"] = "json", - headless: bool = True, - user_agent: str = "", - respect_robots: bool = True, - rate_limit_seconds: float = 1.0, - screenshot: bool = False, - ) -> dict[str, Any]: - """Scrape authorized public/user-owned pages and save structured output.""" - creds = ctx.llm # Read platform-forwarded LLM metadata; this skill does not construct a model. - await ctx.emit_progress(f"starting compliant scrape; llm credential source={creds.source}") - - validation = _validate_inputs( - url=url, - selectors=selectors, - max_pages=max_pages, - max_pages_limit=self.config.max_pages_limit, - wait_seconds=wait_seconds, - output_format=output_format, - user_agent=user_agent, - rate_limit_seconds=rate_limit_seconds, - ) - if validation["errors"]: - return { - "status": "invalid_input", - "items": [], - "saved_paths": [], - "warnings": validation["errors"], - "screenshots": [], - } - - start_url = validation["url"] - effective_user_agent = user_agent.strip() or self.config.default_user_agent - warnings: list[str] = list(validation["warnings"]) - - if respect_robots: - robots = _robots_allowed(start_url, effective_user_agent) - warnings.extend(robots.get("warnings", [])) - if not robots.get("allowed", False): - return { - "status": "blocked_by_robots_txt", - "items": [], - "saved_paths": [], - "warnings": warnings + ["robots.txt disallows fetching the requested URL for this user agent."], - "screenshots": [], - "robots": robots, - } - - script_payload = { - "url": start_url, - "selectors": selectors, - "max_pages": max_pages, - "wait_seconds": wait_seconds, - "output_format": output_format, - "headless": headless, - "user_agent": effective_user_agent, - "respect_robots": respect_robots, - "rate_limit_seconds": rate_limit_seconds, - "screenshot": screenshot, - "output_dir": f"/workspace/{OUTPUT_DIR}", - } - - script = _render_scraper_script(script_payload) - try: - result = await ctx.workspace_python( - script, - image="seleniumbase/seleniumbase:latest", - timeout_seconds=900, - memory_mib=2048, - cpus=2, - ) - except Exception as exc: # noqa: BLE001 - await ctx.emit_error(str(exc), code="browser_runtime_unavailable") - return { - "status": "browser_runtime_unavailable", - "items": [], - "saved_paths": [], - "warnings": warnings - + [ - "Browser sandbox could not be started. The platform must allow the SeleniumBase image, " - "or rerun where browser automation is available.", - f"runtime error: {type(exc).__name__}: {exc}", - ], - "screenshots": [], - } - - stdout = (result.stdout or "").strip() - stderr = (result.stderr or "").strip() - parsed = _parse_last_json(stdout) - if not parsed: - return { - "status": "error", - "items": [], - "saved_paths": [], - "warnings": warnings - + [ - "Browser script did not return a parseable JSON result.", - f"returncode={getattr(result, 'exit_code', None)}", - f"stdout_tail={stdout[-2000:]}", - f"stderr_tail={stderr[-2000:]}", - ], - "screenshots": [], - } - - parsed_warnings = list(parsed.get("warnings") or []) - if getattr(result, "exit_code", 0) not in (0, None): - parsed_warnings.append( - f"browser command returned nonzero exit code {getattr(result, 'exit_code', None)}; usable emitted files were preserved when present" - ) - if stderr: - parsed_warnings.append(f"stderr_tail={stderr[-2000:]}") - - saved_paths = [str(path).replace("/workspace/", "") for path in (parsed.get("saved_paths") or [])] - screenshot_paths = [str(path).replace("/workspace/", "") for path in (parsed.get("screenshots") or [])] - - for path in saved_paths[:5] + screenshot_paths[:5]: - await _emit_workspace_artifact_if_available(ctx, path) - - status = str(parsed.get("status") or "ok") - if status in {"challenge_detected", "manual_intervention_required", "blocked"}: - await ctx.emit_error("Scrape stopped because a CAPTCHA, anti-bot, login, paywall, or access-control challenge was detected.", code=status) - else: - await ctx.emit_progress(f"scrape finished with status={status}; items={len(parsed.get('items') or [])}") - - return { - "status": status, - "items": parsed.get("items") or [], - "saved_paths": saved_paths, - "warnings": warnings + parsed_warnings, - "screenshots": screenshot_paths, - "pages_visited": parsed.get("pages_visited") or [], - "html_captures": [str(path).replace("/workspace/", "") for path in (parsed.get("html_captures") or [])], - "robots": {"checked": bool(respect_robots)}, - } - - @skill( - name="iterate_scrape_goal", - description=( - "Iteratively plan and run compliant scraping passes until a goal is met, " - "using a platform-provided LLM with strict safety rails. The deep agent may " - "refine CSS selectors and pagination hints within the same origin only, and will stop " - "on any CAPTCHA/anti-bot/login/paywall/access-control indicator." - ), - timeout_seconds=900, - cost_class="llm-medium", - grant_outputs_prefix=OUTPUT_DIR, - grant_write_prefixes=(OUTPUT_DIR,), - grant_run_timeout_seconds=900, - ) - async def iterative_scrape_goal( - self, - ctx: RunContext[NoAuth], - goal: str, - start_url: str, - allowed_domains: list[str] | None = None, - seed_selectors: dict[str, str] | None = None, - max_steps: int = 2, - max_pages: int = 1, - min_items: int = 1, - max_runtime_seconds: int = 600, - wait_seconds: float = 3.0, - output_format: Literal["json", "csv"] = "json", - headless: bool = True, - user_agent: str = "", - respect_robots: bool = True, - rate_limit_seconds: float = 1.0, - screenshot: bool = False, - ) -> dict[str, Any]: - """Use an LLM planner to refine selectors/pagination until the goal is met or a stop condition occurs.""" - start_url = str(start_url or "").strip() - if not start_url: - return {"status": "invalid_input", "items": [], "saved_paths": [], "warnings": ["url is required"], "screenshots": []} - - # Initial validation (if selectors are provided) - selectors = seed_selectors or {"title": "title"} - if selectors: - v = _validate_inputs( - url=start_url, - selectors=selectors, - max_pages=max_pages, - max_pages_limit=self.config.max_pages_limit, - wait_seconds=wait_seconds, - output_format=output_format, - user_agent=user_agent, - rate_limit_seconds=rate_limit_seconds, - ) - if v["errors"]: - return {"status": "invalid_input", "items": [], "saved_paths": [], "warnings": v["errors"], "screenshots": []} - - creds = ctx.llm - same_origin_safeguard = _same_origin - await ctx.emit_progress("starting deep-agent iterative scrape planner") - - # Domain guard - def _allowed_origin(url: str) -> bool: - if not allowed_domains: - return same_origin_safeguard(start_url, url) - host = urllib.parse.urlparse(url).netloc.lower() - return any(host.endswith(dom.lower()) for dom in allowed_domains) - - # Planner system prompt with strict safety policy - system_prompt = ( - "You are a cautious scraping planner for a compliant browser scraper. " - "Your actions are limited to proposing: (a) CSS selectors to extract fields; " - "(b) a same-origin URL (relative or absolute) for the next pass; (c) a reason to stop.\n" - "Hard safety rules:\n" - "- Never suggest evading CAPTCHAs, anti-bot, or access controls.\n" - "- Never suggest login, credential use, or paywall bypass.\n" - "- Always respect robots.txt and same-origin restriction with the initial URL's scheme+host.\n" - "- Do NOT use stealth, proxies, or automation-evasion.\n" - "Output only compact JSON with keys: action ('scrape'|'stop'), reason, selectors (object, optional), url (string, optional), confidence (0..1)." - ) - - decisions: list[dict[str, Any]] = [] - step_logs: list[dict[str, Any]] = [] - all_items: list[dict[str, Any]] = [] - all_pages: list[str] = [] - all_htmls: list[str] = [] - all_screens: list[str] = [] - warnings: list[str] = [] - saved_paths: list[str] = [] - - if not creds.api_key: - # Deterministic fallback: run one pass with provided selectors (or trivial title extraction) - await ctx.emit_progress("LLM creds unavailable; running deterministic single-pass fallback") - sel = selectors or {"title": "title"} - one = await self.scrape_website( - ctx=ctx, - url=start_url, - selectors=sel, - max_pages=max_pages, - wait_seconds=wait_seconds, - output_format=output_format, - headless=headless, - user_agent=user_agent, - respect_robots=respect_robots, - rate_limit_seconds=rate_limit_seconds, - screenshot=screenshot, - ) - return { - "status": one.get("status") or "ok", - "items": one.get("items") or [], - "saved_paths": one.get("saved_paths") or [], - "warnings": one.get("warnings") or [], - "screenshots": one.get("screenshots") or [], - "pages_visited": one.get("pages_visited") or [], - "html_captures": one.get("html_captures") or [], - "iterations": 1, - "decisions": [], - } - - # Build deep agent planner - try: - planner = self._build_deep_agent(ctx=ctx, creds=creds, system_prompt=system_prompt) - except Exception as exc: # noqa: BLE001 - warnings.append("planner_unavailable: " + _friendly_error(exc)) - planner = None - current_url = start_url - current_selectors = selectors or {"title": "title"} - - hard_cap_steps = max(1, min(int(max_steps), 3)) - hard_cap_runtime = max(30, min(int(max_runtime_seconds), 840)) - started_at = datetime.now(timezone.utc) - - for step in range(1, hard_cap_steps + 1): - # Runtime budget check - if (datetime.now(timezone.utc) - started_at).total_seconds() > hard_cap_runtime: - warnings.append("max_runtime_seconds reached; stopping") - break - # Execute one browser pass - plan_note = "initial pass" if step == 1 else "refined pass" - run = await self.scrape_website( - ctx=ctx, - url=current_url, - selectors=current_selectors, - max_pages=max_pages, - wait_seconds=wait_seconds, - output_format=output_format, - headless=headless, - user_agent=user_agent, - respect_robots=respect_robots, - rate_limit_seconds=rate_limit_seconds, - screenshot=screenshot, - ) - status = str(run.get("status") or "ok") - items = list(run.get("items") or []) - all_items.extend(items) - all_pages.extend(run.get("pages_visited") or []) - all_htmls.extend(run.get("html_captures") or []) - all_screens.extend(run.get("screenshots") or []) - warnings.extend(run.get("warnings") or []) - saved_paths.extend(run.get("saved_paths") or []) - - step_logs.append({ - "plan": plan_note, - "action": {"url": current_url, "selectors": current_selectors, "max_pages": max_pages}, - "observations": {"status": status, "items_this_pass": len(items), "pages": run.get("pages_visited") or []}, - "safety": {"robots": bool(respect_robots)}, - "decision": None, - }) - if status in {"challenge_detected", "manual_intervention_required", "blocked", "blocked_by_robots_txt", "browser_runtime_unavailable"}: - warnings.append(f"step {step} stopped due to status={status}") - step_logs[-1]["decision"] = {"stop": True, "reason": status} - break - if len(all_items) >= int(min_items): - decisions.append({"step": step, "action": "stop", "reason": "min_items reached", "confidence": 1.0}) - step_logs[-1]["decision"] = {"stop": True, "reason": "min_items reached"} - break - - # Ask planner what to do next - summary = { - "step": step, - "goal": str(goal or "").strip(), - "start_url": start_url, - "current_url": current_url, - "same_origin": same_origin_safeguard(start_url, current_url), - "items_this_pass": len(items), - "items_total": len(all_items), - "example_item": items[0] if items else None, - "current_selectors": current_selectors, - "status": status, - } - prompt = json.dumps( - { - "instruction": "Propose next step using only same-origin URL and safe CSS selectors, or stop.", - "schema": { - "action": "scrape|stop", - "reason": "string", - "selectors": {"field": "CSS"}, - "url": "relative or absolute same-origin URL", - "confidence": 0.0, - }, - "context": summary, - }, - ensure_ascii=False, - ) - if planner is None: - decisions.append({"step": step, "action": "stop", "reason": "planner unavailable", "confidence": 0.0}) - break - state = await planner.ainvoke( - {"messages": [{"role": "user", "content": prompt}]}, - config={"recursion_limit": DEEPAGENTS_RECURSION_LIMIT}, - ) - raw = _last_message_text(state) - try: - decision = _extract_json_object(raw) - except Exception as exc: # noqa: BLE001 - warnings.append("planner response was not JSON; stopping: " + _friendly_error(exc)) - break - - action = str(decision.get("action") or "").strip().lower() - reason = str(decision.get("reason") or "").strip() - next_url = str(decision.get("url") or current_url).strip() or current_url - proposed_selectors = decision.get("selectors") or {} - confidence = float(decision.get("confidence") or 0.0) - - # Enforce safety: allowed_domains and same-origin only - if not _allowed_origin(next_url): - decisions.append({ - "step": step, - "action": "stop", - "reason": "planner suggested URL outside allowed domains; stopping", - "confidence": confidence, - }) - step_logs[-1]["decision"] = {"stop": True, "reason": "outside allowed domains"} - break - - if action == "stop": - decisions.append({"step": step, "action": "stop", "reason": reason, "confidence": confidence}) - step_logs[-1]["decision"] = {"stop": True, "reason": reason} - break - - if action != "scrape": - decisions.append({ - "step": step, - "action": "stop", - "reason": "planner returned unknown action", - "confidence": confidence, - }) - step_logs[-1]["decision"] = {"stop": True, "reason": "unknown action"} - break - - # Adopt next plan - current_url = next_url - if isinstance(proposed_selectors, dict) and proposed_selectors: - safe_map: dict[str, str] = {} - for k, v in proposed_selectors.items(): - k2 = str(k).strip() - v2 = str(v).strip() - if k2 and v2: - safe_map[k2] = v2 - if safe_map: - current_selectors = safe_map - decisions.append({ - "step": step, - "action": "scrape", - "reason": reason or "refining selectors/pagination", - "confidence": confidence, - "url": current_url, - "selectors": current_selectors, - }) - step_logs[-1]["decision"] = {"stop": False, "next_url": current_url, "selectors": current_selectors} - - # Deduplicate items by URL+page_index if present - seen_keys: set[tuple[Any, Any]] = set() - unique_items: list[dict[str, Any]] = [] - for it in all_items: - key = (it.get("url"), it.get("page_index")) - if key in seen_keys: - continue - seen_keys.add(key) - unique_items.append(it) - - # Persist a small run manifest for traceability - manifest_path = await _write_iter_manifest(ctx, { - "goal": goal, - "start_url": start_url, - "timestamp": datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ"), - "steps_count": len(decisions) or min(int(max_steps), 1), - "decisions": decisions, - "steps": step_logs, - "items_total": len(unique_items), - }) - if manifest_path: - saved_paths.append(manifest_path) - - return { - "status": "ok" if unique_items else "no_items", - "items": unique_items, - "saved_paths": saved_paths, - "warnings": warnings, - "screenshots": all_screens, - "pages_visited": all_pages, - "html_captures": all_htmls, - "iterations": len(decisions) or 0, - "steps": len(decisions) or 0, - "decisions": decisions, - } - - def _build_deep_agent(self, *, ctx: RunContext[NoAuth], creds: Any, system_prompt: str): - from a2a_pack.deepagents import create_a2a_deep_agent - from langchain.agents.middleware import wrap_model_call - - @wrap_model_call - async def safe_model_call(request: Any, handler: Any) -> Any: - messages = request.state.get("messages", []) - print( - "[seleniumbase-website-scraper] planner model_call " - f"model={creds.model} source={creds.source} messages={len(messages)}" - ) - return await handler(request) - - backend = ctx.workspace_backend() - skill_sources = _seed_runtime_skills(backend, ctx) - return create_a2a_deep_agent( - ctx, - creds=creds, - backend=backend, - skills=skill_sources or None, - middleware=[safe_model_call], - system_prompt=system_prompt, - ) - - -def _validate_inputs( - *, - url: str, - selectors: dict[str, str], - max_pages: int, - max_pages_limit: int, - wait_seconds: float, - output_format: str, - user_agent: str, - rate_limit_seconds: float, -) -> dict[str, Any]: - errors: list[str] = [] - warnings: list[str] = [] - clean_url = str(url or "").strip() - parsed = urllib.parse.urlparse(clean_url) - if parsed.scheme not in {"http", "https"} or not parsed.netloc: - errors.append("url must be an absolute http(s) URL.") - if parsed.username or parsed.password: - errors.append("url must not contain embedded credentials.") - if not selectors: - errors.append("selectors mapping must contain at least one CSS selector.") - for key, value in selectors.items(): - if not isinstance(key, str) or not key.strip(): - errors.append("selector field names must be non-empty strings.") - if not isinstance(value, str) or not value.strip(): - errors.append(f"selector for {key!r} must be a non-empty string.") - if re.search(r"password|token|secret|credential", key, re.I): - warnings.append(f"selector field {key!r} looks sensitive; do not scrape secrets or credentials.") - if max_pages < 1 or max_pages > max_pages_limit: - errors.append(f"max_pages must be between 1 and {max_pages_limit}.") - if wait_seconds < 0 or wait_seconds > 60: - errors.append("wait_seconds must be between 0 and 60.") - if output_format not in {"json", "csv"}: - errors.append("output_format must be 'json' or 'csv'.") - if "\n" in user_agent or "\r" in user_agent: - errors.append("user_agent must be a single-line string.") - if rate_limit_seconds < 0 or rate_limit_seconds > 120: - errors.append("rate_limit_seconds must be between 0 and 120.") - return {"errors": errors, "warnings": warnings, "url": clean_url} - - -def _robots_allowed(url: str, user_agent: str) -> dict[str, Any]: - parsed = urllib.parse.urlparse(url) - robots_url = urllib.parse.urlunparse((parsed.scheme, parsed.netloc, "/robots.txt", "", "", "")) - rp = urllib.robotparser.RobotFileParser() - rp.set_url(robots_url) - try: - with urllib.request.urlopen( - urllib.request.Request(robots_url, headers={"User-Agent": user_agent}), - timeout=10, - ) as response: - body = response.read(2000000).decode("utf-8", errors="ignore") - rp.parse(body.splitlines()) - return {"allowed": bool(rp.can_fetch(user_agent, url)), "robots_url": robots_url, "warnings": []} - except urllib.error.HTTPError as exc: - if exc.code in {401, 403}: - return { - "allowed": False, - "robots_url": robots_url, - "warnings": [f"robots.txt returned HTTP {exc.code}; treating as disallow for safety."], - } - if exc.code == 404: - return {"allowed": True, "robots_url": robots_url, "warnings": ["robots.txt not found; proceeding because no robots rules were published."]} - return { - "allowed": False, - "robots_url": robots_url, - "warnings": [f"robots.txt check failed with HTTP {exc.code}; treating as disallow for safety."], - } - except Exception as exc: # noqa: BLE001 - return { - "allowed": False, - "robots_url": robots_url, - "warnings": [f"robots.txt could not be checked ({type(exc).__name__}: {exc}); treating as disallow for safety."], - } - - -def _render_scraper_script(payload: dict[str, Any]) -> str: - payload_json = json.dumps(payload, ensure_ascii=False) - return textwrap.dedent( - f""" - import csv - import json - import os - import re - import sys - import time - import traceback - import urllib.parse - import urllib.request - import urllib.robotparser - from datetime import datetime, timezone - - PAYLOAD = json.loads({payload_json!r}) - CHALLENGE_PATTERNS = [ - r"captcha", r"recaptcha", r"hcaptcha", r"cf-challenge", r"cloudflare", r"turnstile", - r"are you human", r"verify you are human", r"bot detection", r"automated traffic", - r"access denied", r"temporarily blocked", r"unusual traffic", r"login required", - r"sign in to continue", r"subscribe to continue", r"paywall", r"forbidden", - ] - PAGINATION_KEYS = {{"next", "_next", "__next__", "pagination_next", "next_page"}} - - def emit(payload): - print(json.dumps(payload, ensure_ascii=False)) - - def safe_name(value): - value = re.sub(r"[^a-zA-Z0-9._-]+", "-", value).strip("-._") - return value[:80] or "page" - - def same_origin(a, b): - pa, pb = urllib.parse.urlparse(a), urllib.parse.urlparse(b) - return pa.scheme == pb.scheme and pa.netloc.lower() == pb.netloc.lower() - - def robots_allowed(url, user_agent): - parsed = urllib.parse.urlparse(url) - robots_url = urllib.parse.urlunparse((parsed.scheme, parsed.netloc, "/robots.txt", "", "", "")) - rp = urllib.robotparser.RobotFileParser() - rp.set_url(robots_url) - try: - req = urllib.request.Request(robots_url, headers={{"User-Agent": user_agent}}) - with urllib.request.urlopen(req, timeout=10) as response: - body = response.read(2000000).decode("utf-8", errors="ignore") - rp.parse(body.splitlines()) - return bool(rp.can_fetch(user_agent, url)), None - except Exception as exc: - return False, "robots.txt could not be checked for pagination URL {{0}}: {{1}}: {{2}}".format(url, type(exc).__name__, exc) - - def split_selector(selector): - selector = selector.strip() - m = re.search(r"::attr\\(([^)]+)\\)\\s*$", selector) - if m: - return selector[:m.start()].strip(), "attr", m.group(1).strip() - if selector.endswith("::text"): - return selector[:-6].strip(), "text", None - if selector.endswith("::html"): - return selector[:-6].strip(), "html", None - return selector, "text", None - - def detect_challenge(page_source, title, current_url): - haystack = "\\n".join([title or "", current_url or "", page_source[:50000] or ""]).lower() - matched = [pat for pat in CHALLENGE_PATTERNS if re.search(pat, haystack, re.I)] - if matched: - return "manual_intervention_required", "challenge/access-control indicator detected: " + ", ".join(sorted(set(matched))[:8]) - return None, None - - def extract_with_bs4(page_source, selectors): - from bs4 import BeautifulSoup - soup = BeautifulSoup(page_source, "html.parser") - item = {{}} - next_url = None - for key, raw_selector in selectors.items(): - css, mode, attr = split_selector(str(raw_selector)) - try: - elements = soup.select(css) - except Exception as exc: - item[key] = None - item.setdefault("_selector_warnings", []).append("selector {{0}}={{1!r}} failed: {{2}}".format(key, raw_selector, exc)) - continue - values = [] - for el in elements: - if mode == "attr": - values.append(el.get(attr)) - elif mode == "html": - values.append(el.decode_contents()) - else: - values.append(el.get_text(" ", strip=True)) - values = [v for v in values if v is not None] - if key in PAGINATION_KEYS: - next_url = values[0] if values else None - else: - item[key] = values if len(values) != 1 else values[0] if values else None - return item, next_url - - def write_outputs(items, pages_visited, html_captures, screenshots, warnings, status): - out_dir = PAYLOAD["output_dir"] - os.makedirs(out_dir, exist_ok=True) - stamp = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") - base = safe_name(urllib.parse.urlparse(PAYLOAD["url"]).netloc) + "-" + stamp - metadata = {{ - "status": status, - "items": items, - "pages_visited": pages_visited, - "html_captures": html_captures, - "screenshots": screenshots, - "warnings": warnings, - "safety": "Authorized public/user-owned scraping only. No CAPTCHA, anti-bot, login, paywall, robots.txt, ToS, or access-control circumvention attempted.", - }} - json_path = os.path.join(out_dir, base + ".json") - with open(json_path, "w", encoding="utf-8") as f: - json.dump(metadata, f, ensure_ascii=False, indent=2) - saved_paths = [json_path] - if PAYLOAD["output_format"] == "csv": - csv_path = os.path.join(out_dir, base + ".csv") - fieldnames = sorted({{k for item in items for k in item.keys() if k != "_selector_warnings"}}) - with open(csv_path, "w", newline="", encoding="utf-8") as f: - writer = csv.DictWriter(f, fieldnames=fieldnames or ["url"]) - writer.writeheader() - for item in items: - row = {{k: (json.dumps(v, ensure_ascii=False) if isinstance(v, (list, dict)) else v) for k, v in item.items() if k in (fieldnames or ["url"])}} - writer.writerow(row) - saved_paths.append(csv_path) - return saved_paths - - def run(): - warnings = [] - items = [] - pages_visited = [] - html_captures = [] - screenshots = [] - status = "ok" - out_dir = PAYLOAD["output_dir"] - os.makedirs(out_dir, exist_ok=True) - current = PAYLOAD["url"] - start = current - - try: - from seleniumbase import SB - browser_kwargs = {{ - "headless": bool(PAYLOAD["headless"]), - "uc": False, - "test": False, - "locale_code": "en", - }} - if PAYLOAD.get("user_agent"): - browser_kwargs["agent"] = PAYLOAD["user_agent"] - with SB(**browser_kwargs) as sb: - for page_index in range(int(PAYLOAD["max_pages"])): - if not same_origin(start, current): - warnings.append("same-domain safeguard stopped pagination to " + str(current)) - break - if PAYLOAD.get("respect_robots"): - allowed, robot_warning = robots_allowed(current, PAYLOAD.get("user_agent") or "*") - if robot_warning: - warnings.append(robot_warning) - if not allowed: - status = "blocked_by_robots_txt" - warnings.append("robots.txt disallows pagination URL " + str(current)) - break - if page_index > 0 and float(PAYLOAD["rate_limit_seconds"]) > 0: - time.sleep(float(PAYLOAD["rate_limit_seconds"])) - - sb.open(current) - if float(PAYLOAD["wait_seconds"]) > 0: - sb.sleep(float(PAYLOAD["wait_seconds"])) - current_url = sb.get_current_url() - title = sb.get_title() - page_source = sb.get_page_source() - pages_visited.append(current_url) - - challenge_status, challenge_warning = detect_challenge(page_source, title, current_url) - if challenge_status: - status = challenge_status - warnings.append(challenge_warning) - html_path = os.path.join(out_dir, "page-{{0}}-challenge.html".format(page_index + 1)) - with open(html_path, "w", encoding="utf-8") as f: - f.write(page_source) - html_captures.append(html_path) - if PAYLOAD.get("screenshot"): - shot_path = os.path.join(out_dir, "page-{{0}}-challenge.png".format(page_index + 1)) - try: - sb.save_screenshot(shot_path) - screenshots.append(shot_path) - except Exception as exc: - warnings.append("screenshot failed: " + str(exc)) - break - - html_path = os.path.join(out_dir, "page-{{0}}.html".format(page_index + 1)) - with open(html_path, "w", encoding="utf-8") as f: - f.write(page_source) - html_captures.append(html_path) - if PAYLOAD.get("screenshot"): - shot_path = os.path.join(out_dir, "page-{{0}}.png".format(page_index + 1)) - try: - sb.save_screenshot(shot_path) - screenshots.append(shot_path) - except Exception as exc: - warnings.append("screenshot failed: " + str(exc)) - - item, next_raw = extract_with_bs4(page_source, PAYLOAD["selectors"]) - item["url"] = current_url - item["page_index"] = page_index + 1 - if item.get("_selector_warnings"): - warnings.extend(item.pop("_selector_warnings")) - items.append(item) - - if page_index >= int(PAYLOAD["max_pages"]) - 1 or not next_raw: - break - next_url = urllib.parse.urljoin(current_url, str(next_raw)) - if not same_origin(start, next_url): - warnings.append("same-domain safeguard refused pagination URL " + str(next_url)) - break - current = next_url - except Exception as exc: - status = "error" - warnings.append("browser automation failed: {{0}}: {{1}}".format(type(exc).__name__, exc)) - warnings.append(traceback.format_exc()[-4000:]) - - saved_paths = write_outputs(items, pages_visited, html_captures, screenshots, warnings, status) - emit({{ - "status": status, - "items": items, - "saved_paths": saved_paths, - "screenshots": screenshots, - "html_captures": html_captures, - "warnings": warnings, - "pages_visited": pages_visited, - }}) - return 0 if status in {{"ok", "blocked_by_robots_txt", "manual_intervention_required"}} else 1 - - if __name__ == "__main__": - sys.exit(run()) - """ - ) - - -def _parse_last_json(stdout: str) -> dict[str, Any] | None: - for line in reversed(stdout.splitlines()): - line = line.strip() - if not line.startswith("{"): - continue - try: - value = json.loads(line) - except json.JSONDecodeError: - continue - if isinstance(value, dict): - return value - return None - - -def _same_origin(a: str, b: str) -> bool: - pa, pb = urllib.parse.urlparse(a), urllib.parse.urlparse(b) - return pa.scheme == pb.scheme and pa.netloc.lower() == pb.netloc.lower() - - -async def _emit_workspace_artifact_if_available(ctx: RunContext[Any], path: str) -> None: - try: - reader = getattr(ctx.workspace, "read_bytes", None) - if reader is None: - return - data = reader(path) - if not isinstance(data, (bytes, bytearray)): - return - mime = "application/json" - if path.endswith(".csv"): - mime = "text/csv" - elif path.endswith(".png"): - mime = "image/png" - elif path.endswith(".html"): - mime = "text/html" - ref = await ctx.write_artifact(Path(path).name, bytes(data), mime) - await ctx.emit_artifact(ref) - except Exception: # noqa: BLE001 - return - - -def _runtime_skills_root(ctx: RunContext[Any]) -> str: - workspace = getattr(ctx, "_workspace", None) - prefixes = tuple(getattr(workspace, "write_prefixes", ()) or ()) - if not prefixes: - outputs_prefix = getattr(workspace, "outputs_prefix", None) - prefixes = (outputs_prefix or "outputs/",) - prefix = str(prefixes[0]).strip("/") - return f"/{prefix}/{RUNTIME_SKILLS_DIR}" if prefix else f"/{RUNTIME_SKILLS_DIR}" - - -def _seed_runtime_skills(backend: Any, ctx: RunContext[Any]) -> list[str]: - # Seed any packaged runtime skills if present - root = Path(__file__).parent / "skills" - if not root.exists(): - return [] - runtime_root = _runtime_skills_root(ctx) - uploads: list[tuple[str, bytes]] = [] - for path in root.rglob("*"): - if path.is_file(): - rel = path.relative_to(root).as_posix() - uploads.append((runtime_root + rel, path.read_bytes())) - if uploads: - backend.upload_files(uploads) - return [runtime_root] - return [] - - -def _extract_json_object(raw: str) -> dict[str, Any]: - text = str(raw or "").strip() - if text.startswith("```"): - text = re.sub(r"^```(?:json)?\s*", "", text) - text = re.sub(r"\s*```$", "", text) - try: - parsed = json.loads(text) - except json.JSONDecodeError: - start = text.find("{") - end = text.rfind("}") - if start < 0 or end <= start: - raise - parsed = json.loads(text[start : end + 1]) - if not isinstance(parsed, dict): - raise ValueError("planner result was not a JSON object") - return parsed - - -def _last_message_text(state: dict[str, Any]) -> str: - messages = state.get("messages") or [] - if not messages: - return json.dumps(state, default=str) - content = getattr(messages[-1], "content", None) - if isinstance(content, str): - return content - if isinstance(content, list): - parts: list[str] = [] - for item in content: - if isinstance(item, dict): - text = item.get("text") or item.get("content") - if text: - parts.append(str(text)) - elif item: - parts.append(str(item)) - return "\n".join(parts) - return str(content or messages[-1]) - - -def _friendly_error(exc: Exception) -> str: - message = str(exc) or exc.__class__.__name__ - message = re.sub(r"(?i)(password|passwd|pwd|secret|token|key)=\S+", r"\1=", message) - message = re.sub(r"(?i)(auth(?:entication)? failed).*", r"\1", message) - return message[:400] - - -async def _write_iter_manifest(ctx: RunContext[Any], manifest: dict[str, Any]) -> str | "": - try: - out_dir = OUTPUT_DIR - name = f"iter-plan-" + datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") + ".json" - path = f"{out_dir}/{name}" - ctx.workspace.write_bytes(path, json.dumps(manifest, ensure_ascii=False, indent=2).encode("utf-8")) - return path - except Exception: # noqa: BLE001 - return ""