From abff09dc0c118627719f63499c4f2d1e9820bd65 Mon Sep 17 00:00:00 2001 From: Navid Filsaraee Date: Mon, 15 Jun 2026 19:42:03 +0330 Subject: [PATCH] init --- .env.example | 22 ++ .gitignore | 6 + admin/__init__.py | 0 admin/main.py | 28 ++ admin/models.py | 241 +++++++++++++ admin/routes.py | 135 +++++++ admin/static/dashboard.html | 558 +++++++++++++++++++++++++++++ config.py | 38 ++ crawler/__init__.py | 0 crawler/driver.py | 83 +++++ crawler/dynamic_flow.py | 282 +++++++++++++++ crawler/flow_runner.py | 104 ++++++ crawler/scenarios/__init__.py | 0 crawler/scenarios/fresh_session.py | 52 +++ crawler/scenarios/otp_login.py | 106 ++++++ main.py | 121 +++++++ pyproject.toml | 24 ++ 17 files changed, 1800 insertions(+) create mode 100644 .env.example create mode 100644 .gitignore create mode 100644 admin/__init__.py create mode 100644 admin/main.py create mode 100644 admin/models.py create mode 100644 admin/routes.py create mode 100644 admin/static/dashboard.html create mode 100644 config.py create mode 100644 crawler/__init__.py create mode 100644 crawler/driver.py create mode 100644 crawler/dynamic_flow.py create mode 100644 crawler/flow_runner.py create mode 100644 crawler/scenarios/__init__.py create mode 100644 crawler/scenarios/fresh_session.py create mode 100644 crawler/scenarios/otp_login.py create mode 100644 main.py create mode 100644 pyproject.toml diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..378928f --- /dev/null +++ b/.env.example @@ -0,0 +1,22 @@ +# Copy to .env and fill in + +TARGET_URL=https://example.com + +# Admin panel credentials +ADMIN_USERNAME=admin +ADMIN_PASSWORD=changeme +ADMIN_HOST=127.0.0.1 +ADMIN_PORT=8000 + +# SQLite DB path +DB_PATH=data/tracker.db + +# Chrome: set to false to watch the browser window +HEADLESS=true +# CHROME_BINARY=/usr/bin/chromium # optional: explicit Chrome path + +# Scenario 1 — comma-separated phone numbers +PHONE_NUMBERS=+989100000001,+989100000002,+989100000003 + +# Flow steps — comma-separated list matching keys in STEP_REGISTRY +FLOW_STEPS=step_home,step_browse diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..1111c9a --- /dev/null +++ b/.gitignore @@ -0,0 +1,6 @@ +.venv/ +__pycache__/ +*.pyc +.env +data/ +*.db diff --git a/admin/__init__.py b/admin/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/admin/main.py b/admin/main.py new file mode 100644 index 0000000..0b35b59 --- /dev/null +++ b/admin/main.py @@ -0,0 +1,28 @@ +"""FastAPI admin application.""" +from __future__ import annotations + +from pathlib import Path + +from fastapi import FastAPI +from fastapi.responses import FileResponse +from fastapi.staticfiles import StaticFiles + +from admin.models import init_db +from admin.routes import router + +app = FastAPI(title="Seed Admin", docs_url=None, redoc_url=None) + +app.include_router(router) + +_STATIC = Path(__file__).parent / "static" +app.mount("/static", StaticFiles(directory=str(_STATIC)), name="static") + + +@app.on_event("startup") +def on_startup() -> None: + init_db() + + +@app.get("/") +def dashboard() -> FileResponse: + return FileResponse(str(_STATIC / "dashboard.html")) diff --git a/admin/models.py b/admin/models.py new file mode 100644 index 0000000..522927d --- /dev/null +++ b/admin/models.py @@ -0,0 +1,241 @@ +"""SQLite schema and DB helpers (aiosqlite).""" +from __future__ import annotations + +import json +import os +import sqlite3 +from pathlib import Path +from typing import Any + +DB_PATH = os.getenv("DB_PATH", "data/tracker.db") + +CREATE_DDL = """ +CREATE TABLE IF NOT EXISTS runs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + scenario TEXT NOT NULL, + identifier TEXT NOT NULL, + success INTEGER NOT NULL, + failure_reason TEXT, + steps_json TEXT NOT NULL DEFAULT '[]', + total_ms INTEGER NOT NULL DEFAULT 0 +); + +CREATE INDEX IF NOT EXISTS idx_runs_scenario ON runs(scenario); +CREATE INDEX IF NOT EXISTS idx_runs_created_at ON runs(created_at); +CREATE INDEX IF NOT EXISTS idx_runs_identifier ON runs(identifier); + +-- Dynamic flow configurations editable from the admin panel +CREATE TABLE IF NOT EXISTS flow_configs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + name TEXT NOT NULL, + is_active INTEGER NOT NULL DEFAULT 0, + + -- Search step + search_texts_json TEXT NOT NULL DEFAULT '[]', -- JSON string[] + search_box_selector TEXT NOT NULL DEFAULT '', + search_box_selector_type TEXT NOT NULL DEFAULT 'css', -- 'css' | 'xpath' + search_submit_selector TEXT NOT NULL DEFAULT '', -- optional — leave blank to use Enter key + search_submit_selector_type TEXT NOT NULL DEFAULT 'css', + search_results_selector TEXT NOT NULL DEFAULT '', -- container that appears after results load + + -- Scroll step + scroll_container_selector TEXT NOT NULL DEFAULT '', -- element to scroll; blank = window + max_scrolls INTEGER NOT NULL DEFAULT 30, + scroll_pause_ms INTEGER NOT NULL DEFAULT 1200, + no_new_content_timeout_ms INTEGER NOT NULL DEFAULT 3000, -- bail if DOM height unchanged after this + + -- Item selection + item_selector_template TEXT NOT NULL DEFAULT '', + -- Template may contain {item_id} placeholder, e.g. "[data-id='{item_id}']" + -- or a plain selector like ".result-card" (picks first match) + item_selector_type TEXT NOT NULL DEFAULT 'css', + target_item_ids_json TEXT NOT NULL DEFAULT '[]' -- JSON string[] — one click per id +); +""" + + +def init_db() -> None: + """Create the DB file and tables (sync, called at startup).""" + path = Path(DB_PATH) + path.parent.mkdir(parents=True, exist_ok=True) + conn = sqlite3.connect(str(path)) + conn.executescript(CREATE_DDL) + conn.commit() + conn.close() + + +# --------------------------------------------------------------------------- +# Async helpers (used by FastAPI routes via aiosqlite) +# --------------------------------------------------------------------------- + +async def insert_run( + scenario: str, + identifier: str, + success: bool, + failure_reason: str | None, + steps: list[dict[str, object]], + total_ms: int, +) -> int: + import aiosqlite + async with aiosqlite.connect(DB_PATH) as db: + cur = await db.execute( + """INSERT INTO runs (scenario, identifier, success, failure_reason, steps_json, total_ms) + VALUES (?, ?, ?, ?, ?, ?)""", + (scenario, identifier, int(success), failure_reason, json.dumps(steps), total_ms), + ) + await db.commit() + return int(cur.lastrowid or 0) + + +# --------------------------------------------------------------------------- +# flow_configs CRUD +# --------------------------------------------------------------------------- + +FlowConfigRow = dict[str, Any] + + +async def list_flow_configs() -> list[FlowConfigRow]: + import aiosqlite + async with aiosqlite.connect(DB_PATH) as db: + db.row_factory = aiosqlite.Row + rows = await (await db.execute( + "SELECT * FROM flow_configs ORDER BY id DESC" + )).fetchall() + return [_parse_flow_row(dict(r)) for r in rows] + + +async def get_flow_config(cfg_id: int) -> FlowConfigRow | None: + import aiosqlite + async with aiosqlite.connect(DB_PATH) as db: + db.row_factory = aiosqlite.Row + row = await (await db.execute( + "SELECT * FROM flow_configs WHERE id=?", (cfg_id,) + )).fetchone() + return _parse_flow_row(dict(row)) if row else None + + +async def get_active_flow_config() -> FlowConfigRow | None: + import aiosqlite + async with aiosqlite.connect(DB_PATH) as db: + db.row_factory = aiosqlite.Row + row = await (await db.execute( + "SELECT * FROM flow_configs WHERE is_active=1 ORDER BY id DESC LIMIT 1" + )).fetchone() + return _parse_flow_row(dict(row)) if row else None + + +async def upsert_flow_config(data: dict[str, Any], cfg_id: int | None = None) -> int: + import aiosqlite + + fields = [ + "name", "is_active", + "search_texts_json", "search_box_selector", "search_box_selector_type", + "search_submit_selector", "search_submit_selector_type", "search_results_selector", + "scroll_container_selector", "max_scrolls", "scroll_pause_ms", "no_new_content_timeout_ms", + "item_selector_template", "item_selector_type", "target_item_ids_json", + ] + + # Serialize list fields + for key in ("search_texts_json", "target_item_ids_json"): + if key in data and isinstance(data[key], list): + data[key] = json.dumps(data[key]) + + async with aiosqlite.connect(DB_PATH) as db: + if data.get("is_active"): + await db.execute("UPDATE flow_configs SET is_active=0") + + if cfg_id is None: + placeholders = ", ".join("?" for _ in fields) + cols = ", ".join(fields) + values = [data.get(f) for f in fields] + cur = await db.execute( + f"INSERT INTO flow_configs ({cols}) VALUES ({placeholders})", values + ) + row_id = int(cur.lastrowid or 0) + else: + set_clause = ", ".join(f"{f}=?" for f in fields) + values_u = [data.get(f) for f in fields] + [cfg_id] + await db.execute(f"UPDATE flow_configs SET {set_clause} WHERE id=?", values_u) + row_id = cfg_id + + await db.commit() + return row_id + + +async def delete_flow_config(cfg_id: int) -> None: + import aiosqlite + async with aiosqlite.connect(DB_PATH) as db: + await db.execute("DELETE FROM flow_configs WHERE id=?", (cfg_id,)) + await db.commit() + + +async def activate_flow_config(cfg_id: int) -> None: + import aiosqlite + async with aiosqlite.connect(DB_PATH) as db: + await db.execute("UPDATE flow_configs SET is_active=0") + await db.execute("UPDATE flow_configs SET is_active=1 WHERE id=?", (cfg_id,)) + await db.commit() + + +def _parse_flow_row(row: dict[str, Any]) -> FlowConfigRow: + for key in ("search_texts_json", "target_item_ids_json"): + if key in row and isinstance(row[key], str): + try: + row[key] = json.loads(row[key]) + except (json.JSONDecodeError, TypeError): + row[key] = [] + return row + + +async def get_stats() -> dict[str, object]: + import aiosqlite + async with aiosqlite.connect(DB_PATH) as db: + db.row_factory = aiosqlite.Row + + total = (await (await db.execute("SELECT COUNT(*) FROM runs")).fetchone())[0] + unique_ids = (await (await db.execute("SELECT COUNT(DISTINCT identifier) FROM runs")).fetchone())[0] + successes = (await (await db.execute("SELECT COUNT(*) FROM runs WHERE success=1")).fetchone())[0] + failures = (await (await db.execute("SELECT COUNT(*) FROM runs WHERE success=0")).fetchone())[0] + + failure_reasons_raw = await ( + await db.execute( + """SELECT failure_reason, COUNT(*) as cnt + FROM runs WHERE success=0 AND failure_reason IS NOT NULL + GROUP BY failure_reason ORDER BY cnt DESC LIMIT 10""" + ) + ).fetchall() + failure_reasons = [{"reason": r["failure_reason"], "count": r["cnt"]} for r in failure_reasons_raw] + + scenario_breakdown_raw = await ( + await db.execute( + """SELECT scenario, COUNT(*) as total, + SUM(success) as ok, + COUNT(*) - SUM(success) as fail + FROM runs GROUP BY scenario""" + ) + ).fetchall() + scenario_breakdown = [ + {"scenario": r["scenario"], "total": r["total"], "success": r["ok"], "failure": r["fail"]} + for r in scenario_breakdown_raw + ] + + recent_raw = await ( + await db.execute( + """SELECT id, created_at, scenario, identifier, success, failure_reason, total_ms + FROM runs ORDER BY id DESC LIMIT 50""" + ) + ).fetchall() + recent = [dict(r) for r in recent_raw] + + return { + "total_runs": total, + "unique_identifiers": unique_ids, + "successes": successes, + "failures": failures, + "success_rate": round(successes / total * 100, 1) if total else 0.0, + "failure_reasons": failure_reasons, + "scenario_breakdown": scenario_breakdown, + "recent_runs": recent, + } diff --git a/admin/routes.py b/admin/routes.py new file mode 100644 index 0000000..248b19b --- /dev/null +++ b/admin/routes.py @@ -0,0 +1,135 @@ +"""FastAPI routes — dashboard API + auth.""" +from __future__ import annotations + +import secrets +from typing import Annotated, Any + +from fastapi import APIRouter, Depends, HTTPException, status +from fastapi.security import HTTPBasic, HTTPBasicCredentials +from pydantic import BaseModel + +from admin.models import ( + activate_flow_config, + delete_flow_config, + get_flow_config, + get_stats, + list_flow_configs, + upsert_flow_config, +) +from config import config + +router = APIRouter() +security = HTTPBasic() + + +def require_auth(credentials: Annotated[HTTPBasicCredentials, Depends(security)]) -> str: + ok_user = secrets.compare_digest(credentials.username.encode(), config.admin_username.encode()) + ok_pass = secrets.compare_digest(credentials.password.encode(), config.admin_password.encode()) + if not (ok_user and ok_pass): + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail="Invalid credentials", + headers={"WWW-Authenticate": "Basic"}, + ) + return credentials.username + + +Auth = Annotated[str, Depends(require_auth)] + + +# --------------------------------------------------------------------------- +# Stats +# --------------------------------------------------------------------------- + +@router.get("/api/stats") +async def stats(_: Auth) -> dict[str, object]: + return await get_stats() + + +# --------------------------------------------------------------------------- +# Flow configs +# --------------------------------------------------------------------------- + +class FlowConfigIn(BaseModel): + name: str + is_active: bool = False + search_texts: list[str] = [] + search_box_selector: str = "" + search_box_selector_type: str = "css" + search_submit_selector: str = "" + search_submit_selector_type: str = "css" + search_results_selector: str = "" + scroll_container_selector: str = "" + max_scrolls: int = 30 + scroll_pause_ms: int = 1200 + no_new_content_timeout_ms: int = 3000 + item_selector_template: str = "" + item_selector_type: str = "css" + target_item_ids: list[str] = [] + + +def _to_db_dict(body: FlowConfigIn) -> dict[str, Any]: + return { + "name": body.name, + "is_active": int(body.is_active), + "search_texts_json": body.search_texts, + "search_box_selector": body.search_box_selector, + "search_box_selector_type": body.search_box_selector_type, + "search_submit_selector": body.search_submit_selector, + "search_submit_selector_type": body.search_submit_selector_type, + "search_results_selector": body.search_results_selector, + "scroll_container_selector": body.scroll_container_selector, + "max_scrolls": body.max_scrolls, + "scroll_pause_ms": body.scroll_pause_ms, + "no_new_content_timeout_ms": body.no_new_content_timeout_ms, + "item_selector_template": body.item_selector_template, + "item_selector_type": body.item_selector_type, + "target_item_ids_json": body.target_item_ids, + } + + +@router.get("/api/flow-configs") +async def list_configs(_: Auth) -> list[dict[str, object]]: + return await list_flow_configs() # type: ignore[return-value] + + +@router.get("/api/flow-configs/{cfg_id}") +async def get_config(cfg_id: int, _: Auth) -> dict[str, object]: + row = await get_flow_config(cfg_id) + if not row: + raise HTTPException(status_code=404, detail="Not found") + return row # type: ignore[return-value] + + +@router.post("/api/flow-configs", status_code=201) +async def create_config(body: FlowConfigIn, _: Auth) -> dict[str, object]: + new_id = await upsert_flow_config(_to_db_dict(body)) + row = await get_flow_config(new_id) + return row or {} # type: ignore[return-value] + + +@router.put("/api/flow-configs/{cfg_id}") +async def update_config(cfg_id: int, body: FlowConfigIn, _: Auth) -> dict[str, object]: + existing = await get_flow_config(cfg_id) + if not existing: + raise HTTPException(status_code=404, detail="Not found") + await upsert_flow_config(_to_db_dict(body), cfg_id=cfg_id) + row = await get_flow_config(cfg_id) + return row or {} # type: ignore[return-value] + + +@router.post("/api/flow-configs/{cfg_id}/activate", status_code=200) +async def activate_config(cfg_id: int, _: Auth) -> dict[str, object]: + existing = await get_flow_config(cfg_id) + if not existing: + raise HTTPException(status_code=404, detail="Not found") + await activate_flow_config(cfg_id) + return {"activated": cfg_id} + + +@router.delete("/api/flow-configs/{cfg_id}", status_code=204) +async def delete_config(cfg_id: int, _: Auth) -> None: + existing = await get_flow_config(cfg_id) + if not existing: + raise HTTPException(status_code=404, detail="Not found") + await delete_flow_config(cfg_id) diff --git a/admin/static/dashboard.html b/admin/static/dashboard.html new file mode 100644 index 0000000..5dfcfaf --- /dev/null +++ b/admin/static/dashboard.html @@ -0,0 +1,558 @@ + + + + + + Seed — Admin + + + + + +
+ +
+ + + + + +
+
+

Seed

LIVE +
+
+ + +
+
+ +
+
Dashboard
+
Flow Configs
+
+ + +
+
+
+
Total Runs
+
Unique Users
+
Successes
+
Failures
+
Success Rate
+
+ +
+
+
Scenario Breakdown
+

No data yet.

+
+
+
Top Failure Reasons
+

No failures recorded.

+
+
+ +
+
Recent Runs
+
+ + + +
#TimeScenarioIdentifierStatusDurationFailure Reason
No runs yet.
+
+
+
+
+ + +
+
+
+
+
Flow Configurations
+
Only one config is active at a time. The crawler picks it up on every run.
+
+ +
+ +
+

No configs yet. Create one to get started.

+
+
+
+ + + + diff --git a/config.py b/config.py new file mode 100644 index 0000000..8706848 --- /dev/null +++ b/config.py @@ -0,0 +1,38 @@ +from __future__ import annotations + +import os +from dataclasses import dataclass, field +from dotenv import load_dotenv + +load_dotenv() + + +@dataclass +class Config: + target_url: str = os.getenv("TARGET_URL", "https://example.com") + admin_username: str = os.getenv("ADMIN_USERNAME", "admin") + admin_password: str = os.getenv("ADMIN_PASSWORD", "changeme") + admin_host: str = os.getenv("ADMIN_HOST", "127.0.0.1") + admin_port: int = int(os.getenv("ADMIN_PORT", "8000")) + db_path: str = os.getenv("DB_PATH", "data/tracker.db") + headless: bool = os.getenv("HEADLESS", "true").lower() == "true" + chrome_binary: str | None = os.getenv("CHROME_BINARY") + + # Scenario 1 — SMS-OTP: list of phone numbers (one per line in env or comma-separated) + phone_numbers: list[str] = field(default_factory=list) + + # Flow steps — ordered list of step names to execute after login + # Override in .env as: FLOW_STEPS=step_home,step_browse,step_checkout + flow_steps: list[str] = field(default_factory=lambda: ["step_home"]) + + def __post_init__(self) -> None: + raw_phones = os.getenv("PHONE_NUMBERS", "") + if raw_phones: + self.phone_numbers = [p.strip() for p in raw_phones.split(",") if p.strip()] + + raw_steps = os.getenv("FLOW_STEPS", "") + if raw_steps: + self.flow_steps = [s.strip() for s in raw_steps.split(",") if s.strip()] + + +config = Config() diff --git a/crawler/__init__.py b/crawler/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/crawler/driver.py b/crawler/driver.py new file mode 100644 index 0000000..dfc953b --- /dev/null +++ b/crawler/driver.py @@ -0,0 +1,83 @@ +"""Stealth browser factory — bypasses Cloudflare/ArvanCloud bot detection.""" +from __future__ import annotations + +import random +import time +from contextlib import contextmanager +from typing import Generator + +import undetected_chromedriver as uc +from selenium.webdriver.chrome.options import Options +from selenium_stealth import stealth + +from config import config + +_USER_AGENTS = [ + "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36", + "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36", + "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36", +] + + +def _build_options(user_agent: str, headless: bool) -> Options: + opts = Options() + if headless: + opts.add_argument("--headless=new") + opts.add_argument(f"--user-agent={user_agent}") + opts.add_argument("--no-sandbox") + opts.add_argument("--disable-dev-shm-usage") + opts.add_argument("--disable-blink-features=AutomationControlled") + opts.add_argument("--disable-infobars") + opts.add_argument("--window-size=1920,1080") + opts.add_experimental_option("excludeSwitches", ["enable-automation"]) + opts.add_experimental_option("useAutomationExtension", False) + if config.chrome_binary: + opts.binary_location = config.chrome_binary + return opts + + +def make_driver(headless: bool | None = None) -> uc.Chrome: + """Return a stealthed undetected Chrome instance.""" + use_headless = config.headless if headless is None else headless + ua = random.choice(_USER_AGENTS) + opts = _build_options(ua, use_headless) + + driver = uc.Chrome(options=opts, use_subprocess=True) + + stealth( + driver, + languages=["en-US", "en"], + vendor="Google Inc.", + platform="Win32", + webgl_vendor="Intel Inc.", + renderer="Intel Iris OpenGL Engine", + fix_hairline=True, + ) + + # Mask navigator.webdriver via CDP + driver.execute_cdp_cmd( + "Page.addScriptToEvaluateOnNewDocument", + { + "source": """ + Object.defineProperty(navigator, 'webdriver', {get: () => undefined}); + Object.defineProperty(navigator, 'plugins', {get: () => [1, 2, 3, 4, 5]}); + Object.defineProperty(navigator, 'languages', {get: () => ['en-US', 'en']}); + """ + }, + ) + return driver + + +@contextmanager +def driver_session(headless: bool | None = None) -> Generator[uc.Chrome, None, None]: + """Context manager that guarantees driver.quit() on exit.""" + driver = make_driver(headless) + try: + yield driver + finally: + driver.quit() + + +def human_delay(lo: float = 0.8, hi: float = 2.5) -> None: + """Sleep for a random human-like interval.""" + time.sleep(random.uniform(lo, hi)) diff --git a/crawler/dynamic_flow.py b/crawler/dynamic_flow.py new file mode 100644 index 0000000..63b7aa6 --- /dev/null +++ b/crawler/dynamic_flow.py @@ -0,0 +1,282 @@ +""" +Dynamic flow executor — reads the active flow config from DB at runtime and runs: + 1. search_step : type each search text into the search box + 2. scroll_step : infinite-scroll until target item is visible or max_scrolls reached + 3. click_step : click each target item id +""" +from __future__ import annotations + +import asyncio +import time +from dataclasses import dataclass, field +from typing import Any + +import undetected_chromedriver as uc +from selenium.common.exceptions import NoSuchElementException, TimeoutException +from selenium.webdriver.common.action_chains import ActionChains +from selenium.webdriver.common.by import By +from selenium.webdriver.common.keys import Keys +from selenium.webdriver.remote.webelement import WebElement +from selenium.webdriver.support import expected_conditions as EC +from selenium.webdriver.support.ui import WebDriverWait + +from crawler.driver import human_delay + + +class DynamicFlowError(Exception): + pass + + +@dataclass +class StepResult: + name: str + success: bool + data: dict[str, Any] = field(default_factory=dict) + error: str | None = None + duration_ms: int = 0 + + +@dataclass +class DynamicFlowResult: + config_id: int + config_name: str + steps: list[StepResult] = field(default_factory=list) + + @property + def success(self) -> bool: + return all(s.success for s in self.steps) + + @property + def failure_reason(self) -> str | None: + for s in self.steps: + if not s.success: + return f"{s.name}: {s.error}" + return None + + +def _by(selector_type: str) -> str: + return By.XPATH if selector_type.lower() == "xpath" else By.CSS_SELECTOR + + +def _find(driver: uc.Chrome, selector: str, selector_type: str) -> WebElement: + return driver.find_element(_by(selector_type), selector) + + +def _find_all(driver: uc.Chrome, selector: str, selector_type: str) -> list[WebElement]: + return driver.find_elements(_by(selector_type), selector) + + +def _wait_for( + driver: uc.Chrome, + selector: str, + selector_type: str, + timeout: float = 15.0, +) -> WebElement: + return WebDriverWait(driver, timeout).until( + EC.presence_of_element_located((_by(selector_type), selector)) + ) + + +class DynamicFlow: + def __init__(self, driver: uc.Chrome, cfg: dict[str, Any]) -> None: + self.driver = driver + self.cfg = cfg + + def run(self) -> DynamicFlowResult: + result = DynamicFlowResult( + config_id=self.cfg["id"], + config_name=self.cfg["name"], + ) + + search_texts: list[str] = self.cfg.get("search_texts_json") or [] + item_ids: list[str] = self.cfg.get("target_item_ids_json") or [] + + for text in search_texts: + step = self._run_search(text) + result.steps.append(step) + if not step.success: + return result + human_delay(0.8, 1.5) + + for item_id in item_ids: + scroll_step = self._run_scroll_and_find(item_id) + result.steps.append(scroll_step) + if not scroll_step.success: + continue # try next item_id, don't abort the whole flow + + click_step = self._run_click(item_id, scroll_step.data.get("element")) + result.steps.append(click_step) + human_delay(1.0, 2.5) + + return result + + # ------------------------------------------------------------------ + # Step: type into search box + # ------------------------------------------------------------------ + def _run_search(self, text: str) -> StepResult: + t0 = time.monotonic() + name = f"search:{text}" + try: + sel = self.cfg["search_box_selector"] + sel_type = self.cfg["search_box_selector_type"] + if not sel: + raise DynamicFlowError("search_box_selector is not configured") + + box = _wait_for(self.driver, sel, sel_type) + box.clear() + human_delay(0.3, 0.6) + + # Type character by character for a human feel + for ch in text: + box.send_keys(ch) + time.sleep(0.04) + + human_delay(0.4, 0.8) + + # Submit: explicit button or Enter + submit_sel = self.cfg.get("search_submit_selector", "") + if submit_sel: + submit_type = self.cfg.get("search_submit_selector_type", "css") + btn = _wait_for(self.driver, submit_sel, submit_type, timeout=5.0) + btn.click() + else: + box.send_keys(Keys.RETURN) + + # Wait for results container if configured + results_sel = self.cfg.get("search_results_selector", "") + if results_sel: + _wait_for(self.driver, results_sel, "css", timeout=15.0) + + human_delay(1.0, 2.0) + return StepResult(name=name, success=True, data={"text": text}, + duration_ms=int((time.monotonic() - t0) * 1000)) + + except Exception as exc: + return StepResult(name=name, success=False, error=str(exc), + duration_ms=int((time.monotonic() - t0) * 1000)) + + # ------------------------------------------------------------------ + # Step: scroll (infinite-scroll) until item selector matches + # ------------------------------------------------------------------ + def _run_scroll_and_find(self, item_id: str) -> StepResult: + t0 = time.monotonic() + name = f"scroll_find:{item_id}" + try: + item_sel = self._build_item_selector(item_id) + item_sel_type = self.cfg.get("item_selector_type", "css") + scroll_container = self.cfg.get("scroll_container_selector", "").strip() + max_scrolls: int = int(self.cfg.get("max_scrolls", 30)) + pause_ms: int = int(self.cfg.get("scroll_pause_ms", 1200)) + no_new_ms: int = int(self.cfg.get("no_new_content_timeout_ms", 3000)) + + # Check if already visible before scrolling + found = self._find_item_visible(item_sel, item_sel_type) + if found: + return StepResult(name=name, success=True, + data={"item_id": item_id, "scrolls": 0, "element": found}, + duration_ms=int((time.monotonic() - t0) * 1000)) + + last_height = self._get_scroll_height(scroll_container) + stale_since: float | None = None + + for scroll_num in range(1, max_scrolls + 1): + self._scroll_down(scroll_container) + time.sleep(pause_ms / 1000) + + found = self._find_item_visible(item_sel, item_sel_type) + if found: + return StepResult(name=name, success=True, + data={"item_id": item_id, "scrolls": scroll_num, "element": found}, + duration_ms=int((time.monotonic() - t0) * 1000)) + + new_height = self._get_scroll_height(scroll_container) + if new_height == last_height: + if stale_since is None: + stale_since = time.monotonic() + elif (time.monotonic() - stale_since) * 1000 >= no_new_ms: + raise DynamicFlowError( + f"No new content after {no_new_ms} ms — reached end of page " + f"without finding item '{item_id}'" + ) + else: + stale_since = None + last_height = new_height + + raise DynamicFlowError( + f"Item '{item_id}' not found after {max_scrolls} scrolls" + ) + + except DynamicFlowError: + raise + except Exception as exc: + return StepResult(name=name, success=False, error=str(exc), + duration_ms=int((time.monotonic() - t0) * 1000)) + + # ------------------------------------------------------------------ + # Step: click the found element + # ------------------------------------------------------------------ + def _run_click(self, item_id: str, element: Any) -> StepResult: + t0 = time.monotonic() + name = f"click:{item_id}" + try: + if element is None: + raise DynamicFlowError("No element reference from scroll step") + + el: WebElement = element + # Scroll element into view and click + self.driver.execute_script("arguments[0].scrollIntoView({block:'center'});", el) + human_delay(0.3, 0.7) + try: + el.click() + except Exception: + # Fallback: JS click + self.driver.execute_script("arguments[0].click();", el) + + human_delay(0.8, 1.5) + return StepResult(name=name, success=True, + data={"item_id": item_id, "url_after": self.driver.current_url}, + duration_ms=int((time.monotonic() - t0) * 1000)) + + except Exception as exc: + return StepResult(name=name, success=False, error=str(exc), + duration_ms=int((time.monotonic() - t0) * 1000)) + + # ------------------------------------------------------------------ + # Helpers + # ------------------------------------------------------------------ + def _build_item_selector(self, item_id: str) -> str: + template: str = self.cfg.get("item_selector_template", "") + if not template: + raise DynamicFlowError("item_selector_template is not configured") + return template.replace("{item_id}", item_id) + + def _find_item_visible(self, selector: str, selector_type: str) -> WebElement | None: + try: + els = _find_all(self.driver, selector, selector_type) + for el in els: + if el.is_displayed(): + return el + except NoSuchElementException: + pass + return None + + def _get_scroll_height(self, container_sel: str) -> int: + if container_sel: + try: + el = _find(self.driver, container_sel, "css") + return int(self.driver.execute_script("return arguments[0].scrollHeight", el)) + except Exception: + pass + return int(self.driver.execute_script("return document.body.scrollHeight")) + + def _scroll_down(self, container_sel: str) -> None: + if container_sel: + try: + el = _find(self.driver, container_sel, "css") + self.driver.execute_script( + "arguments[0].scrollTop = arguments[0].scrollHeight", el + ) + return + except Exception: + pass + self.driver.execute_script("window.scrollTo(0, document.body.scrollHeight)") diff --git a/crawler/flow_runner.py b/crawler/flow_runner.py new file mode 100644 index 0000000..e876fe3 --- /dev/null +++ b/crawler/flow_runner.py @@ -0,0 +1,104 @@ +"""Executes a pre-defined sequence of flow steps and records results.""" +from __future__ import annotations + +import time +from dataclasses import dataclass, field +from typing import Callable + +import undetected_chromedriver as uc + +from config import config +from crawler.driver import human_delay + +# A flow step is a callable that receives the driver and returns arbitrary data. +StepFn = Callable[[uc.Chrome], dict[str, object]] + +# Registry — add your custom step functions here. +# Keys must match the step names in config.flow_steps. +STEP_REGISTRY: dict[str, StepFn] = {} + + +def register_step(name: str) -> Callable[[StepFn], StepFn]: + """Decorator to register a function as a named flow step.""" + def decorator(fn: StepFn) -> StepFn: + STEP_REGISTRY[name] = fn + return fn + return decorator + + +@dataclass +class StepResult: + name: str + success: bool + data: dict[str, object] = field(default_factory=dict) + error: str | None = None + duration_ms: int = 0 + + +@dataclass +class FlowResult: + scenario: str + identifier: str # phone number or "fresh" + steps: list[StepResult] = field(default_factory=list) + + @property + def success(self) -> bool: + return all(s.success for s in self.steps) + + @property + def failure_reason(self) -> str | None: + for s in self.steps: + if not s.success: + return f"{s.name}: {s.error}" + return None + + +class FlowRunner: + def __init__(self, driver: uc.Chrome) -> None: + self.driver = driver + + def run(self, scenario: str, identifier: str) -> FlowResult: + result = FlowResult(scenario=scenario, identifier=identifier) + + for step_name in config.flow_steps: + fn = STEP_REGISTRY.get(step_name) + if fn is None: + result.steps.append( + StepResult( + name=step_name, + success=False, + error=f"Step '{step_name}' not found in registry", + ) + ) + break + + t0 = time.monotonic() + try: + data = fn(self.driver) + duration = int((time.monotonic() - t0) * 1000) + result.steps.append(StepResult(name=step_name, success=True, data=data, duration_ms=duration)) + human_delay(0.5, 1.5) + except Exception as exc: + duration = int((time.monotonic() - t0) * 1000) + result.steps.append( + StepResult(name=step_name, success=False, error=str(exc), duration_ms=duration) + ) + break # abort remaining steps on first failure + + return result + + +# --------------------------------------------------------------------------- +# Example steps — replace with your actual flow logic +# --------------------------------------------------------------------------- + +@register_step("step_home") +def step_home(driver: uc.Chrome) -> dict[str, object]: + return {"url": driver.current_url, "title": driver.title} + + +@register_step("step_browse") +def step_browse(driver: uc.Chrome) -> dict[str, object]: + # Example: navigate to a listing page + human_delay(1.0, 2.0) + return {"url": driver.current_url} diff --git a/crawler/scenarios/__init__.py b/crawler/scenarios/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/crawler/scenarios/fresh_session.py b/crawler/scenarios/fresh_session.py new file mode 100644 index 0000000..f74d730 --- /dev/null +++ b/crawler/scenarios/fresh_session.py @@ -0,0 +1,52 @@ +"""Scenario 2 — fresh cookie jar per run so the site treats the visitor as a new user.""" +from __future__ import annotations + +import undetected_chromedriver as uc +from selenium.common.exceptions import TimeoutException +from selenium.webdriver.support.ui import WebDriverWait + +from config import config +from crawler.driver import human_delay, make_driver + + +class FreshSessionScenario: + """ + Opens the target URL in a brand-new browser profile (no stored cookies, + localStorage, or cache) so each run appears as a first-time visitor. + + Because undetected-chromedriver already creates an isolated temp profile + per instance, simply spinning up a new driver is sufficient — but we also + explicitly delete all cookies after loading to be safe. + """ + + def __init__(self, driver: uc.Chrome, timeout: int = 30) -> None: + self.driver = driver + self.wait = WebDriverWait(driver, timeout) + + def run(self) -> dict[str, object]: + """Navigate to the target, clear any residual state, return page info.""" + self.driver.delete_all_cookies() + + # Clear localStorage / sessionStorage via JS + self.driver.execute_script( + "try { window.localStorage.clear(); window.sessionStorage.clear(); } catch(e) {}" + ) + + self.driver.get(config.target_url) + human_delay(2.0, 4.0) + + # Wait for the page to reach a ready state + self.wait.until( + lambda d: d.execute_script("return document.readyState") == "complete" + ) + + title = self.driver.title + current_url = self.driver.current_url + cookies = self.driver.get_cookies() + + return { + "title": title, + "url": current_url, + "cookie_count": len(cookies), + "cookies": cookies, + } diff --git a/crawler/scenarios/otp_login.py b/crawler/scenarios/otp_login.py new file mode 100644 index 0000000..33f7d14 --- /dev/null +++ b/crawler/scenarios/otp_login.py @@ -0,0 +1,106 @@ +"""Scenario 1 — SMS-OTP login with a batch of phone numbers.""" +from __future__ import annotations + +import time +from typing import Callable + +import undetected_chromedriver as uc +from selenium.common.exceptions import TimeoutException +from selenium.webdriver.common.by import By +from selenium.webdriver.support import expected_conditions as EC +from selenium.webdriver.support.ui import WebDriverWait + +from config import config +from crawler.driver import human_delay + + +class OTPLoginError(Exception): + pass + + +class OTPLoginScenario: + """ + Drives an SMS-OTP login flow. + + The caller supplies an `otp_resolver` — a callable that receives the + phone number and returns the OTP string. Wire it up to your SMS + gateway / SIM-management API. + """ + + def __init__( + self, + driver: uc.Chrome, + phone: str, + otp_resolver: Callable[[str], str], + timeout: int = 30, + ) -> None: + self.driver = driver + self.phone = phone + self.otp_resolver = otp_resolver + self.wait = WebDriverWait(driver, timeout) + + # ------------------------------------------------------------------ + # Override these selectors to match the target site's actual HTML. + # ------------------------------------------------------------------ + _PHONE_FIELD_SELECTOR = (By.CSS_SELECTOR, "input[type='tel'], input[name='phone']") + _SEND_OTP_BTN_SELECTOR = (By.CSS_SELECTOR, "button[type='submit']") + _OTP_FIELD_SELECTOR = (By.CSS_SELECTOR, "input[name='otp'], input[autocomplete='one-time-code']") + _CONFIRM_BTN_SELECTOR = (By.CSS_SELECTOR, "button[type='submit']") + _SUCCESS_INDICATOR = (By.CSS_SELECTOR, "[data-testid='home'], .dashboard, #main-content") + + def run(self) -> dict[str, object]: + """Execute the OTP login. Returns a result dict.""" + self.driver.get(config.target_url) + human_delay(1.5, 3.0) + + try: + self._enter_phone() + self._request_otp() + otp = self._resolve_otp() + self._enter_otp(otp) + self._confirm_login() + self._wait_for_success() + except TimeoutException as exc: + raise OTPLoginError(f"Timeout during OTP login for {self.phone}") from exc + + return {"phone": self.phone, "cookies": self.driver.get_cookies()} + + def _enter_phone(self) -> None: + field = self.wait.until(EC.element_to_be_clickable(self._PHONE_FIELD_SELECTOR)) + field.clear() + human_delay(0.3, 0.7) + for char in self.phone: + field.send_keys(char) + time.sleep(0.05) + + def _request_otp(self) -> None: + btn = self.wait.until(EC.element_to_be_clickable(self._SEND_OTP_BTN_SELECTOR)) + human_delay(0.5, 1.2) + btn.click() + + def _resolve_otp(self) -> str: + # Poll for the OTP from the external resolver (SMS gateway / webhook) + for attempt in range(12): + human_delay(5.0, 8.0) + otp = self.otp_resolver(self.phone) + if otp: + return otp + if attempt == 11: + raise OTPLoginError(f"OTP not received for {self.phone} after 12 attempts") + return "" # unreachable but satisfies type checker + + def _enter_otp(self, otp: str) -> None: + field = self.wait.until(EC.element_to_be_clickable(self._OTP_FIELD_SELECTOR)) + field.clear() + human_delay(0.3, 0.7) + for char in otp: + field.send_keys(char) + time.sleep(0.08) + + def _confirm_login(self) -> None: + btn = self.wait.until(EC.element_to_be_clickable(self._CONFIRM_BTN_SELECTOR)) + human_delay(0.5, 1.0) + btn.click() + + def _wait_for_success(self) -> None: + self.wait.until(EC.presence_of_element_located(self._SUCCESS_INDICATOR)) diff --git a/main.py b/main.py new file mode 100644 index 0000000..b14ca17 --- /dev/null +++ b/main.py @@ -0,0 +1,121 @@ +"""Entry point — run crawl scenarios or start the admin panel.""" +from __future__ import annotations + +import argparse +import asyncio +import sys +import time + +import uvicorn + +from admin.models import get_active_flow_config, init_db, insert_run +from config import config +from crawler.driver import driver_session +from crawler.dynamic_flow import DynamicFlow +from crawler.scenarios.fresh_session import FreshSessionScenario +from crawler.scenarios.otp_login import OTPLoginScenario + + +def _load_active_flow() -> dict[str, object]: + cfg = asyncio.run(get_active_flow_config()) + if cfg is None: + print("[WARN] No active flow config — running without dynamic flow steps.", file=sys.stderr) + return cfg or {} + + +def _record( + scenario: str, + identifier: str, + flow_result: object | None, + exc: Exception | None, + t0: float, +) -> None: + from crawler.dynamic_flow import DynamicFlowResult + + total_ms = int((time.monotonic() - t0) * 1000) + if exc is not None: + asyncio.run(insert_run(scenario, identifier, False, str(exc), [], total_ms)) + print(f"[FAIL] {identifier}: {exc}", file=sys.stderr) + return + + if isinstance(flow_result, DynamicFlowResult): + steps_data = [ + {"name": s.name, "success": s.success, "error": s.error, "duration_ms": s.duration_ms} + for s in flow_result.steps + ] + asyncio.run(insert_run( + scenario, identifier, flow_result.success, + flow_result.failure_reason, steps_data, total_ms, + )) + status = "OK" if flow_result.success else f"FAIL ({flow_result.failure_reason})" + else: + asyncio.run(insert_run(scenario, identifier, True, None, [], total_ms)) + status = "OK (no flow config)" + + print(f"[{status}] {identifier} — {total_ms} ms") + + +def run_otp(phone: str) -> None: + def otp_resolver(p: str) -> str: + # TODO: wire up your SMS gateway / webhook here. + raise NotImplementedError("Implement otp_resolver to fetch OTP from your SMS gateway") + + flow_cfg = _load_active_flow() + t0 = time.monotonic() + try: + with driver_session() as driver: + OTPLoginScenario(driver, phone, otp_resolver).run() + result = DynamicFlow(driver, flow_cfg).run() if flow_cfg else None + except Exception as exc: + _record("otp", phone, None, exc, t0) + return + _record("otp", phone, result, None, t0) + + +def run_fresh() -> None: + flow_cfg = _load_active_flow() + t0 = time.monotonic() + try: + with driver_session() as driver: + FreshSessionScenario(driver).run() + result = DynamicFlow(driver, flow_cfg).run() if flow_cfg else None + except Exception as exc: + _record("fresh", "fresh", None, exc, t0) + return + _record("fresh", "fresh", result, None, t0) + + +def run_admin() -> None: + init_db() + from admin.main import app + uvicorn.run(app, host=config.admin_host, port=config.admin_port) + + +def main() -> None: + parser = argparse.ArgumentParser(description="Seed crawler & admin") + sub = parser.add_subparsers(dest="cmd", required=True) + + sub.add_parser("admin", help="Start the admin panel") + + otp_p = sub.add_parser("otp", help="Run Scenario 1 (SMS-OTP login)") + otp_p.add_argument("--phone", help="Single phone number (default: all from config)") + + sub.add_parser("fresh", help="Run Scenario 2 (fresh session)") + + args = parser.parse_args() + + if args.cmd == "admin": + run_admin() + elif args.cmd == "otp": + phones = [args.phone] if args.phone else config.phone_numbers + if not phones: + print("No phone numbers configured. Set PHONE_NUMBERS in .env or pass --phone", file=sys.stderr) + sys.exit(1) + for phone in phones: + run_otp(phone) + elif args.cmd == "fresh": + run_fresh() + + +if __name__ == "__main__": + main() diff --git a/pyproject.toml b/pyproject.toml new file mode 100644 index 0000000..d2563ba --- /dev/null +++ b/pyproject.toml @@ -0,0 +1,24 @@ +[project] +name = "seed" +version = "0.1.0" +requires-python = ">=3.11" +dependencies = [ + "undetected-chromedriver>=3.5.5", + "selenium>=4.18.0", + "selenium-stealth>=1.0.6", + "fastapi>=0.111.0", + "uvicorn[standard]>=0.30.0", + "aiosqlite>=0.20.0", + "httpx>=0.27.0", + "python-multipart>=0.0.9", + "pydantic>=2.7.0", + "python-dotenv>=1.0.1", +] + +[tool.pyright] +pythonVersion = "3.11" +typeCheckingMode = "strict" +reportMissingImports = true +reportMissingTypeStubs = false +venvPath = "." +venv = ".venv"