diff --git "a/app.py" "b/app.py" deleted file mode 100644--- "a/app.py" +++ /dev/null @@ -1,2694 +0,0 @@ -""" -THE Z AI — Computer Mode Server v13 — CDP DOM ELEMENTS + STALENESS GUARD + COMBINED STATE -===================================================================================== -v13 يضيف فوق v12 تحسينات مبنية على نفس فلسفة CDP الموجودة (بدون OmniParser — -CDP يعطي إحداثيات حقيقية 100% من DOM مباشرة، وهو أدق وأخف من نموذج رؤية إضافي). -هذه التحسينات مستوحاة من فلسفة agent loop عند Manus (خطوة واحدة، ملاحظة شاملة، -قرار) وليست نسخة من بنيته التحتية السحابية (microVMs/E2B) التي تتطلب بنية -سحابية على مستوى شركة ولا يمكن تكرارها بمجرد تعديل كود. - - ✅ كل ما هو في v12 (CDP، browser_get_elements، browser_click_element، إلخ) - - 🆕 v13 — إجراءات جديدة: - - browser_type_into_element → يكتب مباشرة في حقل بـ index (click+type - في نداء واحد، مثل Manus browser_input(index, text)) - - browser_find_element → بحث نصي عن عنصر بالوصف بدل تخمين id يدوياً - - browser_get_state → نداء واحد يجمع (عناصر + نص الصفحة + screenshot) - بدل 3 نداءات منفصلة — يقلل عدد الجولات قبل كل قرار فعلي - - staleness guard: click/type تتحقق أن الصفحة لم تتغيّر منذ آخر - browser_get_elements قبل استخدام الإحداثيات المخزّنة (يمنع نقر/كتابة - بإحداثيات قديمة بعد أي navigation لم يُبلَّغ عنه صراحة) - - 📋 كيف يشتغل CDP هنا: - - Chromium يُشغَّل مع --remote-debugging-port=XXXX (منفذ خاص لكل مستخدم) - - WebSocket مباشر من السيرفر لـ Chromium DevTools Protocol - - نقرا DOM.getDocument + DOM.querySelectorAll لكل العناصر القابلة للنقر - - نرجع قائمة JSON منظمة: [{id, tag, text, x, y, w, h, href, type}, ...] - - 🔄 التوافق الخلفي: كل أدوات v11/v12 (mouse_click, keyboard_type, إلخ) تعمل بدون تغيير -""" - -import asyncio -import base64 -import hashlib -import io -import json -import os -import re -import shutil -import subprocess -import sys -import time -import urllib.parse -import threading -import weakref - -from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Query -from fastapi.middleware.cors import CORSMiddleware -from fastapi.responses import JSONResponse, HTMLResponse -import uvicorn - -# ── CDP (Chrome DevTools Protocol) ────────────────────────────── -# websockets مكتبة خفيفة للاتصال بـ Chromium مباشرة عبر CDP -# تُثبَّت مع: pip install websockets -try: - import websockets - _CDP_AVAILABLE = True -except ImportError: - _CDP_AVAILABLE = False - print("⚠️ websockets not installed — CDP features disabled. Run: pip install websockets") - -# ════════════════════════════════════════════════════════════════ -# ── إعداد Display Pool ───────────────────────────────────────── -# ════════════════════════════════════════════════════════════════ - -# بداية نطاق الـ displays الافتراضية (يمكن تغييره) -DISPLAY_BASE = 100 -DISPLAY_MAX = 200 # أقصى عدد مستخدمين متزامنين - -_display_lock = threading.Lock() - -# user_id → display_info -# display_info = { -# "display": ":101", -# "xvfb_proc": subprocess.Popen | None, -# "browser_proc": subprocess.Popen | None, -# "last_bg_shot_ts": float, -# "last_bg_hash": str, -# "active_ws": WebSocket | None, # الاتصال الحالي لهذا المستخدم -# } -_user_displays: dict[str, dict] = {} -_display_numbers: set[int] = set() # أرقام الـ displays المستخدمة - -# semaphore للتحكم في عمليات terminal (global لكل السيرفر) -_terminal_sem = asyncio.Semaphore(8) - -# ── CDP Port Pool ──────────────────────────────────────────────── -# كل مستخدم يحصل على منفذ CDP خاص لـ Chromium (9222، 9223، ...) -# هذا يسمح للسيرفر بالتحدث مع متصفح كل مستخدم مستقلاً تماماً -CDP_PORT_BASE = 9222 -_cdp_port_lock = threading.Lock() - -def _get_cdp_port(display: str) -> int: - """يحوّل رقم display لمنفذ CDP مقابل: :100 → 9222، :101 → 9223""" - try: - n = int(display.lstrip(":")) - return CDP_PORT_BASE + (n - DISPLAY_BASE) - except Exception: - return CDP_PORT_BASE - -# ════════════════════════════════════════════════════════════════ -# ── اكتشاف المتصفح ──────────────────────────────────────────── -# ════════════════════════════════════════════════════════════════ - -def _detect_browser() -> str: - """ - v12: نُفضّل Chromium على Firefox لأنه يدعم CDP بشكل أفضل وأكثر استقراراً. - CDP (Chrome DevTools Protocol) يتيح لنا قراءة DOM مباشرة — وهو قلب ميزة - browser_get_elements الجديدة التي تحل محل شبكة الإحداثيات. - إذا ما كانش Chromium متاحاً، نرجع لـ Firefox مع تعطيل ميزات CDP. - """ - # Chromium أولاً (CDP الكامل) - for c in ["chromium-browser", "chromium", "google-chrome", "google-chrome-stable"]: - r = subprocess.run(["which", c], capture_output=True, text=True) - if r.returncode == 0 and r.stdout.strip(): - print(f"🌐 Browser: {c} (CDP ✅ مدعوم)") - return c - # Firefox كبديل (CDP محدود — ميزات DOM ستعمل عبر xdotool فقط) - for c in ["firefox", "firefox-esr"]: - r = subprocess.run(["which", c], capture_output=True, text=True) - if r.returncode == 0 and r.stdout.strip(): - print(f"🌐 Browser: {c} (CDP ⚠️ غير مدعوم — DOM features disabled)") - return c - return "firefox" - -BROWSER = _detect_browser() -# هل المتصفح الحالي يدعم CDP؟ -BROWSER_SUPPORTS_CDP = any(b in BROWSER for b in ["chromium", "chrome"]) -print(f"🔧 CDP Support: {'✅ enabled' if BROWSER_SUPPORTS_CDP else '⚠️ disabled (Firefox detected)'}") - -# ════════════════════════════════════════════════════════════════ -# ── CDP Engine — قلب ميزة "رؤية DOM" المستوحاة من Manus ──────── -# ════════════════════════════════════════════════════════════════ -# -# كيف يشتغل: -# 1. Chromium يُشغَّل مع --remote-debugging-port=PORT -# 2. نتصل بـ http://localhost:PORT/json/list لنجيب webSocketDebuggerUrl -# 3. نُرسل أوامر CDP عبر WebSocket مباشرة للمتصفح -# 4. نستخدم Runtime.evaluate لتشغيل JavaScript يجمع كل العناصر القابلة للنقر -# 5. النتيجة: قائمة JSON بكل عنصر مع إحداثياته الدقيقة (مأخوذة من getBoundingClientRect) -# -# هذا بالضبط ما يفعله Manus عبر browser_view + OmniParser — -# الفرق أننا نأخذ الإحداثيات مباشرة من DOM (أدق وأسرع) بدل كشف YOLO على صورة - -class CDPSession: - """ - جلسة CDP لمستخدم واحد — تُدير الاتصال بـ Chromium WebSocket. - مصممة لإعادة الاستخدام: اتصال واحد مفتوح طوال الجلسة، يُغلق ويُعاد فتحه - تلقائياً عند الحاجة (مثل إعادة تشغيل المتصفح). - """ - def __init__(self, port: int): - self.port = port - self._ws = None # WebSocket connection - self._msg_id = 0 - self._pending: dict = {} # id → asyncio.Future - self._recv_task = None - self._lock = asyncio.Lock() - - async def connect(self, timeout: float = 5.0) -> bool: - """يتصل بـ Chromium CDP — يُعيد True عند النجاح.""" - if not _CDP_AVAILABLE: - return False - try: - import aiohttp - async with aiohttp.ClientSession() as s: - async with s.get( - f"http://localhost:{self.port}/json/list", - timeout=aiohttp.ClientTimeout(total=timeout) - ) as r: - tabs = await r.json() - except Exception: - # aiohttp غير متاح أو Chromium لم يبدأ بعد — نجرب عبر urllib - try: - import urllib.request - with urllib.request.urlopen( - f"http://localhost:{self.port}/json/list", timeout=timeout - ) as r: - tabs = json.loads(r.read()) - except Exception as e: - print(f"[CDP:{self.port}] ❌ Can't reach Chromium: {e}") - return False - - ws_url = None - for tab in tabs: - if tab.get("type") == "page": - ws_url = tab.get("webSocketDebuggerUrl") - break - - if not ws_url: - print(f"[CDP:{self.port}] ❌ No page tab found") - return False - - try: - import websockets as ws_lib - self._ws = await asyncio.wait_for( - ws_lib.connect(ws_url, ping_interval=None, max_size=50 * 1024 * 1024), - timeout=timeout - ) - self._recv_task = asyncio.create_task(self._recv_loop()) - print(f"[CDP:{self.port}] ✅ Connected to Chromium") - return True - except Exception as e: - print(f"[CDP:{self.port}] ❌ WebSocket failed: {e}") - return False - - async def _recv_loop(self): - """يستقبل رسائل CDP باستمرار ويُوزّعها على الـ futures المنتظِرة.""" - try: - async for raw in self._ws: - try: - msg = json.loads(raw) - msg_id = msg.get("id") - if msg_id and msg_id in self._pending: - self._pending[msg_id].set_result(msg) - except Exception: - pass - except Exception: - pass - finally: - # انتهى الاتصال — نُصفّر كل الـ futures المعلّقة - for fut in self._pending.values(): - if not fut.done(): - fut.set_exception(ConnectionError("CDP disconnected")) - self._pending.clear() - self._ws = None - - async def send(self, method: str, params: dict = None, timeout: float = 10.0): - """يُرسل أمر CDP ويُعيد الرد — None عند الخطأ.""" - if not self._ws: - return None - async with self._lock: - self._msg_id += 1 - msg_id = self._msg_id - - fut = asyncio.get_event_loop().create_future() - self._pending[msg_id] = fut - try: - await self._ws.send(json.dumps({ - "id": msg_id, "method": method, - "params": params or {} - })) - result = await asyncio.wait_for(fut, timeout=timeout) - return result - except Exception as e: - self._pending.pop(msg_id, None) - return None - - async def close(self): - if self._recv_task: - self._recv_task.cancel() - if self._ws: - try: - await self._ws.close() - except Exception: - pass - self._ws = None - - -# map: display → CDPSession -_cdp_sessions: dict[str, CDPSession] = {} -_cdp_sessions_lock = threading.Lock() - - -async def _get_cdp_session(display: str) -> "CDPSession | None": - """يُعيد CDPSession نشطة للـ display، أو None لو مش متاح.""" - if not BROWSER_SUPPORTS_CDP or not _CDP_AVAILABLE: - return None - - port = _get_cdp_port(display) - with _cdp_sessions_lock: - sess = _cdp_sessions.get(display) - - if sess and sess._ws: - return sess - - # إنشاء session جديدة أو إعادة اتصال - new_sess = CDPSession(port) - ok = await new_sess.connect(timeout=3.0) - if ok: - with _cdp_sessions_lock: - _cdp_sessions[display] = new_sess - return new_sess - return None - - -# ── JavaScript يجمع كل العناصر القابلة للنقر من DOM ────────────── -# هذا هو القلب التقني للميزة — يعمل كـ OmniParser لكن بدون GPU: -# بدل نموذج YOLO يكشف العناصر من صورة، نسأل المتصفح مباشرة عن عناصر DOM -# ونحصل على إحداثيات دقيقة 100% (getBoundingClientRect) + وصف وظيفي (text/aria) -_CDP_COLLECT_ELEMENTS_JS = """ -(function() { - // عناصر المحدد: كل ما يمكن النقر عليه أو التفاعل معه - const SELECTORS = [ - 'a[href]', // روابط - 'button', // أزرار - 'input', // حقول إدخال - 'select', // قوائم منسدلة - 'textarea', // مناطق نصية - '[onclick]', // عناصر بها onclick - '[role="button"]', // أزرار ARIA - '[role="link"]', // روابط ARIA - '[role="menuitem"]',// عناصر قائمة - '[role="tab"]', // تبويبات - '[role="checkbox"]',// مربعات اختيار - '[role="radio"]', // زر راديو - '[role="combobox"]',// قوائم مركّبة - '[role="listbox"]', // قوائم - '[role="option"]', // خيارات قائمة - '[tabindex]:not([tabindex="-1"])', // عناصر قابلة للتنقل بـ Tab - ]; - - const seen = new Set(); - const results = []; - let id = 0; - - const allEls = document.querySelectorAll(SELECTORS.join(',')); - - for (const el of allEls) { - // تجاهل العناصر غير المرئية - const style = window.getComputedStyle(el); - if (style.display === 'none' || style.visibility === 'hidden' || - style.opacity === '0' || el.offsetParent === null) continue; - - const rect = el.getBoundingClientRect(); - // تجاهل العناصر خارج نطاق الشاشة أو بمساحة صغيرة جداً - if (rect.width < 2 || rect.height < 2) continue; - if (rect.bottom < 0 || rect.top > window.innerHeight * 2) continue; - if (rect.right < 0 || rect.left > window.innerWidth * 2) continue; - - // مفتاح فريد لمنع التكرار - const key = `${Math.round(rect.x)}_${Math.round(rect.y)}_${el.tagName}`; - if (seen.has(key)) continue; - seen.add(key); - - // استخراج النص الوصفي - let text = ( - el.getAttribute('aria-label') || - el.getAttribute('title') || - el.getAttribute('placeholder') || - el.textContent?.trim() || - el.getAttribute('value') || - el.getAttribute('alt') || - '' - ).substring(0, 120).replace(/\\s+/g, ' ').trim(); - - // استخراج معلومات إضافية - const tag = el.tagName.toLowerCase(); - const type = el.getAttribute('type') || ''; - const href = el.getAttribute('href') || ''; - const role = el.getAttribute('role') || tag; - const isInViewport = ( - rect.top >= 0 && rect.bottom <= window.innerHeight && - rect.left >= 0 && rect.right <= window.innerWidth - ); - - // الإحداثيات المطلقة (مع الـ scroll) - const scrollX = window.scrollX || window.pageXOffset; - const scrollY = window.scrollY || window.pageYOffset; - - results.push({ - id: id++, - tag: tag, - role: role, - text: text, - type: type, - href: href, - // إحداثيات viewport (للنقر المباشر) - x: Math.round(rect.x + rect.width / 2), - y: Math.round(rect.y + rect.height / 2), - x1: Math.round(rect.x), - y1: Math.round(rect.y), - w: Math.round(rect.width), - h: Math.round(rect.height), - // إحداثيات الصفحة الكاملة (مع scroll) - page_x: Math.round(rect.x + scrollX + rect.width / 2), - page_y: Math.round(rect.y + scrollY + rect.height / 2), - in_viewport: isInViewport, - }); - } - - return JSON.stringify({ - elements: results, - page_title: document.title, - page_url: window.location.href, - viewport_w: window.innerWidth, - viewport_h: window.innerHeight, - scroll_x: Math.round(window.scrollX), - scroll_y: Math.round(window.scrollY), - total_height: document.body.scrollHeight, - }); -})(); -""" - - -async def cdp_get_elements(display: str) -> dict: - """ - يجمع كل العناصر القابلة للنقر من DOM الصفحة الحالية. - يُعيد dict فيه: - - elements: قائمة العناصر [{id, tag, text, x, y, w, h, ...}] - - page_title, page_url, viewport_w, viewport_h, scroll_y - - error: رسالة خطأ إذا فشل - """ - sess = await _get_cdp_session(display) - if not sess: - return {"error": "CDP not available (Chromium required)", "elements": []} - - try: - result = await sess.send("Runtime.evaluate", { - "expression": _CDP_COLLECT_ELEMENTS_JS, - "returnByValue": True, - "awaitPromise": False, - "timeout": 8000, - }, timeout=10.0) - - if not result: - return {"error": "CDP timeout or no response", "elements": []} - - val = result.get("result", {}).get("result", {}).get("value") - if not val: - exc = result.get("result", {}).get("exceptionDetails") - err = exc.get("text", "Unknown JS error") if exc else "Empty response" - return {"error": err, "elements": []} - - data = json.loads(val) - return data - - except Exception as e: - return {"error": str(e), "elements": []} - - -async def cdp_click_element(display: str, element_id: int, elements: list) -> dict: - """ - ينقر على عنصر بـ index (مثل Manus browser_click(index=N)). - يستخدم إحداثيات viewport مأخوذة مباشرة من DOM. - يُعيد {"ok": True, "x": ..., "y": ...} أو {"error": "..."} - """ - # ابحث عن العنصر في القائمة - el = next((e for e in elements if e.get("id") == element_id), None) - if not el: - return {"error": f"Element {element_id} not found in elements list"} - - x, y = el["x"], el["y"] - - # إذا العنصر خارج viewport، مرّر الصفحة إليه أولاً - if not el.get("in_viewport", True): - sess = await _get_cdp_session(display) - if sess: - await sess.send("Runtime.evaluate", { - "expression": f"document.querySelectorAll('*')[{element_id}]?.scrollIntoView({{block:'center'}})", - "returnByValue": True, - }) - await asyncio.sleep(0.3) - # أعد احتساب الإحداثيات بعد التمرير - new_data = await cdp_get_elements(display) - new_el = next((e for e in new_data.get("elements", []) if e.get("id") == element_id), None) - if new_el: - x, y = new_el["x"], new_el["y"] - - # النقر بالإحداثيات الحقيقية عبر xdotool (نفس محرك v11) - return {"ok": True, "x": x, "y": y, "tag": el.get("tag"), "text": el.get("text", "")[:50]} - - -async def cdp_type_into_element(display: str, element_id: int, elements: list, text: str, - clear_first: bool = True) -> dict: - """ - يكتب نصاً مباشرة داخل عنصر إدخال بـ index — مثل Manus browser_input(index, text). - يدمج خطوتين كانتا منفصلتين (click ثم keyboard_type) في نداء واحد، ويتحقق - أن العنصر فعلاً حقل إدخال قابل للكتابة قبل المحاولة (بدل الكتابة العمياء - في أي عنصر، وهو خطأ شائع عند تحويل click+type لخطوة واحدة). - يُعيد {"ok": True, "x": ..., "y": ...} أو {"error": "..."}. - """ - el = next((e for e in elements if e.get("id") == element_id), None) - if not el: - return {"error": f"Element {element_id} not found in elements list"} - - tag = (el.get("tag") or "").lower() - role = (el.get("role") or "").lower() - editable_tags = {"input", "textarea"} - editable_roles = {"textbox", "searchbox", "combobox"} - if tag not in editable_tags and role not in editable_roles: - return { - "error": f"Element {element_id} is a {tag or role}, not a text input. " - f"Use browser_click_element if you meant to click it." - } - - x, y = el["x"], el["y"] - if not el.get("in_viewport", True): - new_data = await cdp_get_elements(display) - new_el = next((e for e in new_data.get("elements", []) if e.get("id") == element_id), None) - if new_el: - x, y = new_el["x"], new_el["y"] - - return {"ok": True, "x": x, "y": y, "tag": tag or role, "clear_first": clear_first} - - -async def cdp_get_page_text(display: str) -> dict: - """يجيب نص الصفحة الكامل (document.body.innerText) — مفيد لفهم المحتوى بدون صورة.""" - sess = await _get_cdp_session(display) - if not sess: - return {"error": "CDP not available", "text": ""} - - result = await sess.send("Runtime.evaluate", { - "expression": """(function(){ - return JSON.stringify({ - text: document.body?.innerText?.substring(0, 50000) || '', - title: document.title, - url: window.location.href, - char_count: document.body?.innerText?.length || 0, - }); - })()""", - "returnByValue": True, - "timeout": 5000, - }) - - try: - val = result["result"]["result"]["value"] - return json.loads(val) - except Exception as e: - return {"error": str(e), "text": ""} - - -async def cdp_execute_js(display: str, js_code: str) -> dict: - """يُنفّذ JavaScript مباشرة في المتصفح — مثل Manus browser_console_exec.""" - sess = await _get_cdp_session(display) - if not sess: - return {"error": "CDP not available", "result": None} - - result = await sess.send("Runtime.evaluate", { - "expression": js_code, - "returnByValue": True, - "awaitPromise": True, - "timeout": 10000, - }, timeout=12.0) - - if not result: - return {"error": "CDP timeout", "result": None} - - r = result.get("result", {}).get("result", {}) - exc = result.get("result", {}).get("exceptionDetails") - if exc: - return {"error": exc.get("text", "JS Error"), "result": None} - - return {"result": r.get("value"), "type": r.get("type", "undefined")} - - -async def _elements_are_stale(display: str, cached_url: str) -> bool: - """ - يتحقق أن الصفحة الحالية ما زالت نفسها اللي التُقطت منها last_elements. - إحداثيات عنصر مأخوذة من صفحة قديمة قد تشير الآن لعنصر مختلف تماماً - (أو فراغ) بعد أي navigation — هذا فحص رخيص (URL فقط) يمنع نقرات/كتابة - عمياء في عنصر خاطئ بعد تنقّل لم يُبلَّغ عنه صراحة. - """ - if not cached_url: - return False # لا معلومة سابقة كافية — لا نمنع، فقط لا نتحقق - current = await cdp_get_page_text(display) - current_url = current.get("url", "") - if not current_url: - return False # فشل الفحص نفسه — لا نمنع تنفيذ الأمر بسببه - return current_url != cached_url - - -def cdp_find_elements_by_text(elements: list, query: str, limit: int = 8) -> list: - """ - يبحث عن عناصر بمطابقة نصية جزئية (case-insensitive) ف�� النص/الوصف/href. - يسمح للـ LLM بقول "دوس على زر تسجيل الدخول" بدل قراءة القائمة الكاملة - وتخمين رقم id — أقرب لطريقة تفكير Manus عن الواجهة (بالوصف الوظيفي - مش بمسح يدوي لكل عنصر). - """ - if not query or not elements: - return [] - q = query.strip().lower() - scored = [] - for el in elements: - haystack = " ".join([ - str(el.get("text", "")), - str(el.get("href", "")), - str(el.get("role", "")), - str(el.get("type", "")), - ]).lower() - if q in haystack: - # تفضيل التطابق الأدق (بداية النص) على التطابق الجزئي فقط - score = 0 if str(el.get("text", "")).lower().startswith(q) else 1 - scored.append((score, el)) - scored.sort(key=lambda t: t[0]) - return [e for _, e in scored[:limit]] - - -async def cdp_navigate(display: str, url: str) -> dict: - """يُنقل المتصفح لرابط معين عبر CDP — أسرع وأكثر موثوقية من xdotool.""" - sess = await _get_cdp_session(display) - if not sess: - return {"error": "CDP not available"} - - result = await sess.send("Page.navigate", {"url": url}, timeout=15.0) - if not result: - return {"error": "Navigation timeout"} - - frame_id = result.get("result", {}).get("frameId") - if frame_id: - # انتظر اكتمال تحميل الصفحة - await asyncio.sleep(0.5) - return {"ok": True, "url": url} - return {"error": result.get("result", {}).get("errorText", "Unknown error")} - - -# ── CLI flags ثابتة تُضاف لكل تشغيل Firefox ── -# تسرّع الإقلاع (لا GPU فيزيائي حقيقي في Xvfb فلا داعي لمحاولة -# التسريع المرئي الذي يفشل دائماً هنا ويسبّب تأخيراً في الإقلاع)، -# وتمنع أي nag screens أو crash-reporter windows من الظهور. -FIREFOX_CLI_FLAGS = ["--no-remote", "--new-instance"] -FIREFOX_ENV_EXTRA = { - # يمنع محاولات تسريع GPU غير المتوفرة في Xvfb من إبطاء الإقلاع - "MOZ_DISABLE_GPU_SANDBOX": "1", - "MOZ_ACCELERATED": "0", -} - -# ════════════════════════════════════════════════════════════════ -# ── Firefox Profile (per-user, isolated, crash-recovery DISABLED) ─ -# ════════════════════════════════════════════════════════════════ -# المشكلة الجذرية التي هذا القسم يحلّها: -# عند قتل Firefox (حتى بـ pkill "لطيف")، لا يُغلق بشكل نظيف دائماً على -# سيرفر بموارد محدودة/Xvfb، فتُسجَّل الجلسة على أنها "crashed" داخل -# sessionstore.jsonlz4 الخاص بالبروفايل. في المرة التالية يفتح Firefox -# ويعرض شاشة "Sorry. We're having trouble getting your pages back" بدل -# الصفحة المطلوبة — وهذه هي الشاشة السوداء/الفاشلة التي تصل كلقطة "ناجحة" -# تقنياً (ليست سوداء فعلاً) لكنها لا تحتوي على المحتوى المطلوب أبداً. -# الحل الجذري: بروفايل خاص بكل مستخدم (isolated) مع user.js يعطّل تماماً: -# - استرجاع الجلسة بعد التحطّم (session restore prompt) -# - أي محاولة لإعادة فتح تبويبات سابقة -# بهذا، كل فتح لفايرفوكس يبدأ صفحة نظيفة فارغة دائماً، بغض النظر عن كيفية -# إغلاقه سابقاً. - -FIREFOX_PROFILE_PREFS = """ -// ── تعطيل استرجاع الجلسة نهائياً (سبب شاشة "Restore Session") ── -user_pref("browser.sessionstore.resume_from_crash", false); -user_pref("browser.sessionstore.resume_session_once", false); -user_pref("browser.sessionstore.max_resumed_crashes", -1); -user_pref("browser.sessionstore.restore_on_demand", false); -user_pref("browser.sessionstore.enabled", false); -user_pref("browser.sessionstore.privacy_level", 2); -user_pref("browser.startup.page", 0); -user_pref("browser.startup.homepage_override.mstone", "ignore"); -user_pref("toolkit.startup.max_resumed_crashes", -1); -// تعطيل نافذة "Restore previous session" ونوافذ الأزمات المختلفة -user_pref("browser.sessionstore.max_tabs_undo", 0); -user_pref("browser.sessionstore.max_windows_undo", 0); -user_pref("browser.tabs.crashReporting.sendReport", false); -user_pref("browser.crashReports.unsubmittedCheck.autoSubmit2", false); -user_pref("browser.crashReports.unsubmittedCheck.enabled", false); -user_pref("toolkit.crashreporter.infoURL", ""); -// تعطيل شاشة "Restore" وكذلك أي إشعار تحطّم -user_pref("browser.sessionstore.resumeFromCrash", false); -user_pref("dom.ipc.plugins.flash.subprocess.crashreporter.enabled", false); -// إلغاء استعادة آخر جلسة تماماً + دائماً about:blank -user_pref("browser.startup.homepage", "about:blank"); -// تعطيل تحديثات ورسائل onboarding التي قد تبطئ أو تحجب اللقطة -user_pref("browser.shell.checkDefaultBrowser", false); -user_pref("browser.aboutwelcome.enabled", false); -user_pref("browser.startup.firstrunSkipsHomepage", true); -user_pref("startup.homepage_welcome_url", ""); -user_pref("startup.homepage_welcome_url.additional", ""); -user_pref("browser.uitour.enabled", false); -user_pref("browser.newtabpage.activity-stream.feeds.telemetry", false); -user_pref("datareporting.policy.dataSubmissionEnabled", false); -user_pref("app.normandy.enabled", false); -user_pref("app.update.enabled", false); -user_pref("app.update.auto", false); -// ── إصلاحات تجمد التحميل ("Connecting to ... تبقى للأبد") ── -// السبب الغالب على سيرفرات مجانية/container: محاولة IPv6 أولاً (غير -// مدعوم فعلياً على أغلب مزودي Render/HuggingFace المجانية) قبل الرجوع -// لـ IPv4 بعد timeout طويل جداً (قد يصل 20-30 ثانية لكل طلب DNS)، -// وهذا بالضبط السبب الذي يظهر كشاشة "Connecting..." معلقة لا تتقدم أبداً. -user_pref("network.dns.disableIPv6", true); -user_pref("network.http.fast-fallback-to-IPv4", true); -// تقليل مهلة connect الإجمالية لكل طلب من 15 ثانية (الافتراضي) إلى 8 ثوانٍ -// — أفضل أن يفشل الطلب بسرعة ويقرر الذكاء الاصطناعي بديلاً من أن يبقى عالقاً طويلاً. -user_pref("network.http.connection-timeout", 8); -user_pref("network.http.response.timeout", 12); -// تعطيل predictive prefetching/speculative connect التي تستهلك عرض حزمة على -// سيرفر محدود الموارد بدون فائدة حقيقية هنا. -user_pref("network.dns.disablePrefetch", true); -user_pref("network.prefetch-next", false); -user_pref("network.predictor.enabled", false); -user_pref("network.http.speculative-parallel-limit", 0); -// تعطيل safebrowsing الذي يستدعي طلبات خارجية إضافية عند فتح كل صفحة -// (يُبطئ التحميل أو يعلقه لو الطلب لقوائم Google فشل أو استغرق). -user_pref("browser.safebrowsing.malware.enabled", false); -user_pref("browser.safebrowsing.phishing.enabled", false); -user_pref("browser.safebrowsing.downloads.enabled", false); -user_pref("browser.safebrowsing.provider.google4.updateURL", ""); -user_pref("browser.safebrowsing.provider.google.updateURL", ""); -// تعطيل telemetry/captive-portal checks التي تفتح اتصالات خلفية غير ضرورية عند الإقلاع -user_pref("network.captive-portal-service.enabled", false); -user_pref("network.connectivity-service.enabled", false); -user_pref("toolkit.telemetry.server", ""); -// تقليل عدد اتصالات HTTP المتزامنة لكل دومين — يقلل الضغط على شبكة محدودة -// النطاق للسيرفرات المجانية ويقلل احتمال التعلق. -user_pref("network.http.max-persistent-connections-per-server", 4); -user_pref("network.http.max-connections", 48); -""" - -def _ensure_firefox_profile(user_id: str, display: str) -> str: - """ - يُنشئ (أو يُعيد استخدام) بروفايل Firefox مخصص ونظيف لهذا المستخدم، - مع إعدادات تعطّل استرجاع الجلسة (crash recovery) نهائياً. - يُعيد المسار الكامل للبروفايل. - """ - safe_id = re.sub(r"[^a-zA-Z0-9_.-]", "_", user_id) or "anon" - profile_dir = os.path.expanduser(f"~/.zai_ff_profiles/{safe_id}") - try: - os.makedirs(profile_dir, exist_ok=True) - prefs_path = os.path.join(profile_dir, "user.js") - with open(prefs_path, "w", encoding="utf-8") as f: - f.write(FIREFOX_PROFILE_PREFS) - except Exception as e: - print(f"[profile] ⚠️ failed to prepare profile for {user_id}: {e}") - return profile_dir - - -def _wipe_firefox_session_data(profile_dir: str): - """ - يمسح ملفات الجلسة الفاسدة (sessionstore) + أي lock متبقٍّ. - هذا ضروري لأن مجرد تعطيل session restore في user.js لا يمسح ملفات - جلسة سابقة موجودة بالفعل على القرص من قبل تفعيل هذه الإعدادات. - """ - try: - subprocess.run( - f"rm -f '{profile_dir}'/lock '{profile_dir}'/.parentlock " - f"'{profile_dir}'/sessionstore.jsonlz4 " - f"'{profile_dir}'/sessionstore-backups/*.jsonlz4 " - f"'{profile_dir}'/sessionCheckpoints.json 2>/dev/null", - shell=True, timeout=5, capture_output=True - ) - except Exception: - pass - - -# ════════════════════════════════════════════════════════════════ -# ── Xvfb Management (per-display) ──────────────────────────── -# ════════════════════════════════════════════════════════════════ - -def _next_free_display() -> int: - """يجد رقم display حر غير مستخدم.""" - used = _display_numbers.copy() - for n in range(DISPLAY_BASE, DISPLAY_MAX): - if n not in used: - return n - # إذا امتلأت القائمة — أعد استخدام أقدم display غير نشط - return DISPLAY_BASE - - -def _start_chromium_with_cdp(display: str, url: str = "about:blank") -> "subprocess.Popen | None": - """ - يُشغّل Chromium مع CDP مفعّل على منفذ خاص بهذا الـ display. - v12: هذا يحل محل Firefox ويُضيف إمكانية قراءة DOM مباشرة. - - الـ flags المُستخدمة: - --remote-debugging-port → CDP WebSocket (قلب الميزة الجديدة) - --disable-gpu → ضروري داخل Xvfb (لا GPU حقيقي) - --no-sandbox → مطلوب في بيئات Docker/container - --disable-dev-shm-usage → يمنع مشاكل shared memory في containers - --disable-extensions → أسرع إقلاع - --window-size → حجم ثابت = حجم Xvfb (1280x800) - """ - port = _get_cdp_port(display) - env = {**os.environ, "DISPLAY": display} - - # أوقف أي Chromium قديم على نفس المنفذ - try: - subprocess.run( - ["pkill", "-9", "-f", f"[r]emote-debugging-port={port}"], - capture_output=True, timeout=3 - ) - time.sleep(0.3) - except Exception: - pass - - chromium_cmd = [ - BROWSER, - f"--remote-debugging-port={port}", - "--remote-debugging-address=127.0.0.1", # محلي فقط (أمان) - "--disable-gpu", - "--disable-software-rasterizer", - "--no-sandbox", - "--disable-dev-shm-usage", - "--disable-extensions", - "--disable-notifications", - "--disable-popup-blocking", - "--disable-translate", - "--disable-background-timer-throttling", - "--disable-backgrounding-occluded-windows", - "--disable-renderer-backgrounding", - "--no-first-run", - "--no-default-browser-check", - "--disable-default-apps", - "--disable-sync", - "--disable-logging", - "--silent-debugger-extension-api", - "--window-size=1280,800", - "--window-position=0,0", - f"--user-data-dir=/tmp/chromium-cdp-{port}", - url, - ] - - try: - proc = subprocess.Popen( - chromium_cmd, env=env, - stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL - ) - print(f"[chromium:{display}] 🚀 Started on CDP port {port} (PID {proc.pid})") - return proc - except FileNotFoundError: - print(f"[chromium:{display}] ❌ Chromium binary not found: {BROWSER}") - return None - except Exception as e: - print(f"[chromium:{display}] ❌ Failed to start: {e}") - return None - - -def _start_xvfb(display: str) -> "subprocess.Popen | None": - """يُشغّل Xvfb على display محدد ويُعيد الـ Popen أو None.""" - # تحقق أولاً: هل الـ display يعمل بالفعل؟ - try: - r = subprocess.run( - ["xdpyinfo", "-display", display], - capture_output=True, timeout=3 - ) - if r.returncode == 0: - print(f"[xvfb] ✅ Display {display} already active") - return None # يعمل بالفعل بدون Popen نشغّله - except Exception: - pass - - try: - # ── دقة 1280x800 بدل 1920x1080 ────────────────────────────── - # هذا يقلل حجم بيانات كل لقطة شاشة بنسبة ~60%، وبالتالي يسرّع كل - # مراحل الالتقاط والترميز بشكل كبير على سيرفر بموارد محدودة. - # نفس الدقة المستخدمة في التطبيقات المرجعية لـ computer-use. - proc = subprocess.Popen( - ["Xvfb", display, "-screen", "0", "1280x800x24", - "-nolisten", "tcp", "-ac"], - stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL - ) - # انتظر حتى يصبح جاهزاً (max 6 ثواني) - for _ in range(60): - time.sleep(0.1) - try: - r = subprocess.run( - ["xdpyinfo", "-display", display], - capture_output=True, timeout=2 - ) - if r.returncode == 0: - print(f"[xvfb] ✅ Xvfb started on {display}") - # إصلاح: Xvfb بدون أي رسم لاحق يكون framebuffer فارغاً تماماً - # (لا لون خلفية أصلاً)، فأول لقطة شاشة قد تبدو "فارغة بشكل - # مربك" رغم أنها التقاط ناجح تقنياً لشاشة نظيفة فعلاً. نرسم - # لوناً محايداً بسيطاً فوراً عبر xsetroot لتوضيح أن الشاشة - # جاهزة وفارغة عمداً (سطح مكتب نظيف) لا أنها خطأ في الالتقاط. - try: - env_xs = {**os.environ, "DISPLAY": display} - subprocess.run( - ["xsetroot", "-solid", "#2b2b2b"], - env=env_xs, timeout=2, capture_output=True - ) - except Exception: - pass - return proc - except Exception: - continue - print(f"[xvfb] ⚠️ Xvfb may not be ready yet on {display}") - return proc - except FileNotFoundError: - print(f"[xvfb] ❌ Xvfb binary not found") - return None - except Exception as e: - print(f"[xvfb] ❌ Failed to start on {display}: {e}") - return None - - -def _kill_proc(proc: "subprocess.Popen | None"): - """يقتل process بأمان.""" - if not proc: - return - try: - proc.terminate() - proc.wait(timeout=3) - except Exception: - try: - proc.kill() - except Exception: - pass - - -def _kill_display_processes(display: str): - """يقتل كل العمليات المرتبطة بـ display معين.""" - env_d = {**os.environ, "DISPLAY": display} - # قتل المتصفح — SIGKILL مباشر (-9) لتفادي أي حالة إغلاق نصفي تُسجَّل - # لاحقاً كتحطّم (crash) داخل Firefox وتُنتج شاشة "Restore Session" - for b in ["firefox", "firefox-esr", "chromium", "chrome"]: - try: - subprocess.run( - ["pkill", "-9", "-f", f"[{b[0]}]{b[1:]}.*{display}"], - timeout=3, capture_output=True - ) - except Exception: - pass - # قتل Xvfb على هذا الـ display - try: - subprocess.run( - ["pkill", "-f", f"[X]vfb {display}"], - timeout=3, capture_output=True - ) - except Exception: - pass - time.sleep(0.5) - - -# ════════════════════════════════════════════════════════════════ -# ── Session Management (per user_id) ───────────────────────── -# ════════════════════════════════════════════════════════════════ - -async def get_or_create_user_session(user_id: str, ws: WebSocket) -> dict: - """ - يُعيد session المستخدم (أو يُنشئ واحدة جديدة إذا لم تكن موجودة). - كل user_id ← display خاص + Xvfb خاص. - """ - with _display_lock: - if user_id in _user_displays: - sess = _user_displays[user_id] - # حدّث الـ WebSocket الحالي - sess["active_ws"] = ws - print(f"[session] 🔁 Reconnected user '{user_id}' on display {sess['display']}") - return sess - - # مستخدم جديد — خصص له display - disp_num = _next_free_display() - _display_numbers.add(disp_num) - display = f":{disp_num}" - - sess = { - "user_id": user_id, - "display": display, - "xvfb_proc": None, # سيُشغَّل لاحقاً - "browser_proc": None, - "last_bg_shot_ts": 0.0, - "last_bg_hash": "", - "active_ws": ws, - "created": time.time(), - "consecutive_capture_failures": 0, # عدّاد فشل الالتقاط المتتالي — يُصفَّر عند أي نجاح أو reset - "cdp_port": _get_cdp_port(display), # v12: منفذ CDP الخاص بهذا المستخدم - "last_elements": [], # v12: آخر قائمة عناصر DOM (للـ click_element) - } - _user_displays[user_id] = sess - print(f"[session] ✅ New user '{user_id}' → display {display}") - return sess - - -async def reset_user_computer(user_id: str) -> dict: - """ - يُعيد ضبط الكمبيوتر الافتراضي للمستخدم: - يقتل Xvfb والمتصفح ويُشغّل Xvfb جديداً نظيفاً. - """ - with _display_lock: - sess = _user_displays.get(user_id) - if not sess: - return {} - - display = sess["display"] - print(f"[reset] 🔄 Resetting computer for user '{user_id}' on {display}") - - # قتل المتصفح — SIGKILL مباشر (بدل terminate اللطيف) لأن الأخير قد - # يترك فايرفوكس في حالة نصف-مغلقة تُسجَّل كتحطّم (crash) في المرة - # القادمة، وهذا بالضبط ما يُنتج شاشة "Restore Session" لاحقاً. - browser_proc = sess.get("browser_proc") - if browser_proc: - try: - browser_proc.kill() - browser_proc.wait(timeout=3) - except Exception: - pass - sess["browser_proc"] = None - - # قتل Xvfb - _kill_proc(sess.get("xvfb_proc")) - sess["xvfb_proc"] = None - - # تنظيف شامل لهذا الـ display - _kill_display_processes(display) - - # ── تنظيف شامل لبيانات جلسة Firefox الفاسدة (السبب الجذري لشاشة - # "Restore Session") — ليس فقط lock files بل sessionstore كاملاً، - # في كل من البروفايل الافتراضي وأي بروفايل مخصص لهذا المستخدم ── - try: - subprocess.run( - "rm -f ~/.mozilla/firefox/*/lock ~/.mozilla/firefox/*/.parentlock " - "~/.mozilla/firefox/*/sessionstore.jsonlz4 " - "~/.mozilla/firefox/*/sessionstore-backups/*.jsonlz4 " - "~/.mozilla/firefox/*/sessionCheckpoints.json 2>/dev/null", - shell=True, timeout=5, capture_output=True - ) - except Exception: - pass - safe_id = re.sub(r"[^a-zA-Z0-9_.-]", "_", user_id) or "anon" - custom_profile = os.path.expanduser(f"~/.zai_ff_profiles/{safe_id}") - _wipe_firefox_session_data(custom_profile) - - # تشغيل Xvfb جديد (خارج الـ lock لأن start_xvfb يستغرق وقتاً) - new_proc = await asyncio.to_thread(_start_xvfb, display) - with _display_lock: - if user_id in _user_displays: - _user_displays[user_id]["xvfb_proc"] = new_proc - _user_displays[user_id]["last_bg_shot_ts"] = 0.0 - _user_displays[user_id]["last_bg_hash"] = "" - _user_displays[user_id]["consecutive_capture_failures"] = 0 - - print(f"[reset] ✅ Computer reset done for '{user_id}' on {display}") - return _user_displays.get(user_id, {}) - - -async def ensure_xvfb_for_session(sess: dict): - """يتأكد أن Xvfb يعمل لهذا الـ session — يُشغّله إذا لم يكن كذلك.""" - display = sess["display"] - - # تحقق إذا كان يعمل بالفعل - try: - r = subprocess.run( - ["xdpyinfo", "-display", display], - capture_output=True, timeout=3 - ) - if r.returncode == 0: - return # يعمل - except Exception: - pass - - # شغّله - proc = await asyncio.to_thread(_start_xvfb, display) - with _display_lock: - if sess["user_id"] in _user_displays: - _user_displays[sess["user_id"]]["xvfb_proc"] = proc - - -async def destroy_user_ws(user_id: str, ws: WebSocket): - """ - يُزيل الـ WebSocket من الـ session عند انقطاع الاتصال. - لا يحذف الـ session نفسها — المستخدم يحتفظ بكمبيوتره. - """ - with _display_lock: - sess = _user_displays.get(user_id) - if sess and sess.get("active_ws") is ws: - sess["active_ws"] = None - print(f"[session] 📴 User '{user_id}' disconnected (session kept)") - - -# ════════════════════════════════════════════════════════════════ -# ── Screenshot Engine (per-display) ────────────────────────── -# ════════════════════════════════════════════════════════════════ - -def _is_black_screen(img) -> bool: - try: - small = img.resize((100, 100)) - pixels = list(small.getdata()) - avg = sum(sum(p[:3]) for p in pixels) / (len(pixels) * 3 * 255) - return avg < 0.04 - except Exception: - return False - - -def _load_capture_image(path: str): - from PIL import Image - if not path or not os.path.exists(path): - return None, 0, 0 - if os.path.getsize(path) < 1024: - return None, 0, 0 - try: - img = Image.open(path).convert("RGB") - w, h = img.size - if w < 100 or h < 100: - return None, 0, 0 - return img, w, h - except Exception as ex: - print(f"[cap] load error for {path}: {ex}") - return None, 0, 0 - - -def _capture_via_import_pipe(display: str): - """ - الطريقة الأساسية والوحيدة: `import -window root` مع الإخراج مباشرة - عبر stdout (pipe) — بدون أي ملف مؤقت على القرص. أسرع بشكل ملحوظ من - الكتابة لملف ثم إعادة قراءته، خصوصاً على تخزين شبكي بطيء. - """ - from PIL import Image - env = {**os.environ, "DISPLAY": display} - try: - r = subprocess.run( - ["import", "-window", "root", "-silent", "png:-"], - env=env, capture_output=True - ) - if r.returncode != 0 or not r.stdout or len(r.stdout) < 500: - return None, 0, 0 - img = Image.open(io.BytesIO(r.stdout)).convert("RGB") - w, h = img.size - if w < 100 or h < 100: - return None, 0, 0 - return img, w, h - except FileNotFoundError: - print(f"[cap:{display}] import (ImageMagick) not installed") - return None, 0, 0 - except Exception as e: - print(f"[cap:{display}] import-pipe: {e}") - return None, 0, 0 - - -def _capture_via_xlib_direct(display: str): - """ - خط الدفاع الثاني: قراءة X11 framebuffer مباشرة عبر python-xlib، - بدون subprocess إطلاقاً (أسرع من أي أداة خارجية، ويعمل حتى لو - ImageMagick غير مثبّت على السيرفر). - """ - from PIL import Image - try: - from Xlib import display as Xdisp, X - xd = Xdisp.Display(display) - root = xd.screen().root - geom = root.get_geometry() - w, h = geom.width, geom.height - raw = root.get_image(0, 0, w, h, X.ZPixmap, 0xFFFFFFFF) - img = Image.frombuffer("RGB", (w, h), raw.data, "raw", "BGRX", 0, 1) - xd.close() - return img, w, h - except Exception as e: - print(f"[cap:{display}] xlib-direct: {e}") - return None, 0, 0 - - -def _capture_via_xwd(display: str): - """ - خط الدفاع الثالث — طريقة جديدة تماماً، مستقلة كلياً عن الطريقتين - السابقتين في مسارها البرمجي: `xwd` (X Window Dump) هي أداة أصلية من - حزمة x11-apps الأساسية نفسها (وليست ImageMagick مثل `import`، ولا - مكتبة python-xlib عبر اتصال socket مباشر مثل الطريقة الثانية) — - فإذا كان سبب فشل الطريقتين الأخريين متعلقاً تحديداً بـ ImageMagick - (مفقود/معطوب) أو بمكتبة python-xlib نفسها (تعارض إصدار، اتصال عالق)، - فإن xwd تسلك مساراً مختلفاً بالكامل في النواة قد ينجح حين تفشل الطريقتان - الأخريان. تكتب xwd تنسيق XWD الخام إلى stdout مباشرة (بدون ملف مؤقت على - القرص، بنفس فلسفة import-pipe)، ثم نحوّله عبر PIL. - """ - from PIL import Image - env = {**os.environ, "DISPLAY": display} - try: - r = subprocess.run( - ["xwd", "-root", "-display", display, "-silent"], - env=env, capture_output=True - ) - if r.returncode != 0 or not r.stdout or len(r.stdout) < 500: - return None, 0, 0 - # PIL يدعم قراءة تنسيق XWD مباشرة (plugin XVThumb/XWD مدمج) - img = Image.open(io.BytesIO(r.stdout)) - img = img.convert("RGB") - w, h = img.size - if w < 100 or h < 100: - return None, 0, 0 - return img, w, h - except FileNotFoundError: - print(f"[cap:{display}] xwd not installed") - return None, 0, 0 - except Exception as e: - print(f"[cap:{display}] xwd: {e}") - return None, 0, 0 - - -def _capture_raw(display: str) -> tuple: - """ - يلتقط الشاشة بأسرع طريقة موثوقة: import-pipe أولاً (سريعة جداً)، - ثم xlib-direct، وأخيراً xwd كخط دفاع ثالث جديد كلياً بمسار برمجي - مختلف تماماً عن الطريقتين الأوليين (أداة X11 أصلية أخرى غير - ImageMagick وغير python-xlib) — هذا يضمن أن فشل مصدر واحد (مثل - ImageMagick معطوب أو تعارض في مكتبة xlib) لا يعني فشل الالتقاط - بالكامل، لأن xwd تسلك مساراً مستقلاً تماماً وقد تنجح حين تفشل الأخريان. - لا يوجد أي حد زمني (timeout) على أي محاولة — كل محاولة تُترك تكتمل - فعلياً مهما طال الوقت؛ الحماية من التعليق الفعلي (X server ميت) تأتي - من طبقة أعلى (شمّاعة الفشل المتتالي + reset تلقائي) وليس من قطع - العملية في المنتصف. - - ��─ القيمة الرابعة المُعادة (is_placeholder) ── - True فقط عندما فشلت الطرق الثلاث فعلياً ورجعنا لصورة "Screenshot - failed" المرسومة يدوياً. الطبقة الأعلى (capture_with_grid) تستخدم - هذا العلم لتفعيل عدّاد الفشل المتتالي وتحفيز reset تلقائي كامل — - بدل الاعتماد على فحص محتوى النص داخل الصورة. - """ - for method_name, method in ( - ("import-pipe", _capture_via_import_pipe), - ("xlib-direct", _capture_via_xlib_direct), - ("xwd", _capture_via_xwd), - ): - for attempt in range(2): - try: - img, w, h = method(display) - if img and not _is_black_screen(img): - if attempt > 0: - print(f"[cap:{display}] ✅ {method_name} succeeded on retry") - return img, w, h, False - except Exception as e: - print(f"[cap:{display}] {method_name}: {e}") - if attempt == 0: - time.sleep(0.3) - - print(f"[cap:{display}] ⚠️ All methods failed (import-pipe, xlib-direct, xwd) — placeholder") - from PIL import Image as PILImg, ImageDraw - sw, sh = _get_screen_size(display) - img = PILImg.new("RGB", (sw or 1280, sh or 800), (15, 20, 40)) - draw = ImageDraw.Draw(img) - draw.rectangle([(0, 0), (sw, 55)], fill=(30, 40, 80)) - draw.text((10, 10), f"⚠️ Screenshot failed — DISPLAY={display}", fill=(255, 120, 80)) - draw.text((10, 32), "Methods tried: import-pipe, xlib-direct, xwd", fill=(130, 130, 150)) - return img, sw or 1280, sh or 800, True - - -def _get_screen_size(display: str) -> tuple: - try: - r = subprocess.run( - ["xdotool", "getdisplaygeometry"], - env={**os.environ, "DISPLAY": display}, - capture_output=True, text=True - ) - parts = r.stdout.strip().split() - return int(parts[0]), int(parts[1]) - except Exception: - return 1280, 800 - - -def _get_mouse_pos(display: str) -> tuple: - try: - r = subprocess.run( - ["xdotool", "getmouselocation"], - env={**os.environ, "DISPLAY": display}, - capture_output=True, text=True - ) - mx = int(re.search(r"x:(\d+)", r.stdout).group(1)) - my = int(re.search(r"y:(\d+)", r.stdout).group(1)) - return mx, my - except Exception: - return 0, 0 - - -def _draw_mouse_marker(draw, msx, msy, color=(255, 50, 50, 240)): - """يرسم علامة الماوس (دائرة + خطوط تقاطع) في نقطة محددة على صورة.""" - r = 11 - line_color = (color[0], color[1], color[2], 200) - draw.ellipse([(msx-r, msy-r), (msx+r, msy+r)], outline=color, width=2) - draw.line([(msx-18, msy), (msx+18, msy)], fill=line_color, width=1) - draw.line([(msx, msy-18), (msx, msy+18)], fill=line_color, width=1) - - -def _render_grid_variant(base_img, ow, oh, sw, sh, mx, my, display, - grid_step: int, line_color: tuple, label_every: int, - major_only_labels: bool = True): - """ - يرسم نسخة شبكة إحداثيات واحدة فوق نسخة من الصورة الأساسية. - grid_step → المسافة بالبكسل الحقيقي بين كل خط شبكة. - line_color → لون الخطوط والأرقام (RGBA). - label_every → كل كم بكسل تُكتب فيه تسمية إحداثي (x,y) كاملة عند التقاطعات. - major_only_labels → إذا True، الأرقام على الحواف تظهر فقط عند خطوط "رئيسية". - """ - from PIL import ImageDraw - img = base_img.copy() - draw = ImageDraw.Draw(img, "RGBA") - - step_x = max(1, int(grid_step * sw / ow)) - step_y = max(1, int(grid_step * sh / oh)) - minor_color = (line_color[0], line_color[1], line_color[2], 35) - major_color = (line_color[0], line_color[1], line_color[2], 110) - text_color = line_color - - x_sc, x_r = step_x, grid_step - x_majors = [] - while x_sc < sw: - is_major = (x_r % label_every == 0) - draw.line([(x_sc, 0), (x_sc, sh)], fill=(major_color if is_major else minor_color), width=1) - if is_major or not major_only_labels: - draw.rectangle([(x_sc+1, 2), (x_sc+34, 15)], fill=(0, 0, 0, 175)) - draw.text((x_sc+2, 3), str(x_r), fill=text_color) - if is_major: - x_majors.append((x_sc, x_r)) - x_sc += step_x; x_r += grid_step - - y_sc, y_r = step_y, grid_step - y_majors = [] - while y_sc < sh: - is_major = (y_r % label_every == 0) - draw.line([(0, y_sc), (sw, y_sc)], fill=(major_color if is_major else minor_color), width=1) - if is_major or not major_only_labels: - draw.rectangle([(2, y_sc+1), (38, y_sc+14)], fill=(0, 0, 0, 175)) - draw.text((3, y_sc+2), str(y_r), fill=text_color) - if is_major: - y_majors.append((y_sc, y_r)) - y_sc += step_y; y_r += grid_step - - for (xs, xr) in x_majors: - for (ys, yr) in y_majors: - label = f"{xr},{yr}" - tw = 6 * len(label) + 4 - draw.rectangle([(xs+2, ys+2), (xs+2+tw, ys+13)], fill=(0, 0, 0, 150)) - draw.text((xs+4, ys+2), label, fill=text_color) - - msx = int(mx * sw / ow) - msy = int(my * sh / oh) - _draw_mouse_marker(draw, msx, msy, color=(line_color[0], line_color[1], line_color[2], 240)) - - final = img.convert("RGB") - draw2 = ImageDraw.Draw(final) - draw2.rectangle([(0, 0), (sw, 20)], fill=(0, 0, 0)) - draw2.text((4, 3), f"SCREEN {ow}x{oh} | MOUSE:({mx},{my}) | GRID={grid_step}px | DSP:{display}", fill=(0, 220, 160)) - draw2.rectangle([(0, sh-20), (sw, sh)], fill=(0, 0, 0)) - draw2.text((4, sh-17), "CLICKCOORDS = numbers at every intersection (real screen pixels)", fill=(180, 180, 70)) - return final - - -def capture_with_grid(display: str, scale: float = 0.85, quality: int = 72, - force_mx: int | None = None, - force_my: int | None = None, - grid_step: int = 50) -> dict: - """ - يلتقط الشاشة من display محدد ويُنتج 3 نسخ من نفس اللقطة بالضبط (نفس اللحظة): - - "data" → النسخة العادية (نظيفة تماماً) + علامة الماوس فقط، بدون أي Grid. - - "data_grid" → شبكة إحداثيات عادية (كل 50px)، أخضر/سماوي — للذكاء الاصطناعي. - - "data_grid2" → شبكة إحداثيات أدق (كل 20px)، أحمر — للنقرات الدقيقة. - —— تبسيط v13: نسختان فقط (بدل 3) لتقليل عدد عمليات resize/encode لكل - لقطة بنسبة قرابة س4إلى س2، ودقة المقاس الواحد 1280x800 الجديدة لم - تعد تحتاج لثلاث تدرجات دقة منفصلة لتحديد الإحداثيات. - """ - from PIL import Image, ImageDraw - - img, ow, oh, is_placeholder = _capture_raw(display) - if img is None: - return {"data": "", "data_grid": "", "data_grid2": "", "data_grid3": "", - "width": 1280, "height": 800, "mouse_x": 0, "mouse_y": 0, - "is_placeholder": True} - - mx, my = (force_mx, force_my) if force_mx is not None else _get_mouse_pos(display) - - sw = int(ow * scale) - sh = int(oh * scale) - base_img = img.resize((sw, sh), Image.LANCZOS) - msx = int(mx * sw / ow) - msy = int(my * sh / oh) - - clean_img = base_img.copy() - draw_clean = ImageDraw.Draw(clean_img, "RGBA") - _draw_mouse_marker(draw_clean, msx, msy) - clean_final = clean_img.convert("RGB") - buf_clean = io.BytesIO() - clean_final.save(buf_clean, format="JPEG", quality=quality, optimize=True) - data_clean = base64.b64encode(buf_clean.getvalue()).decode() - - grid1 = _render_grid_variant(base_img, ow, oh, sw, sh, mx, my, display, - grid_step=grid_step, line_color=(0, 255, 180, 235), - label_every=100) - buf1 = io.BytesIO(); grid1.save(buf1, format="JPEG", quality=quality, optimize=True) - data_grid1 = base64.b64encode(buf1.getvalue()).decode() - - grid2 = _render_grid_variant(base_img, ow, oh, sw, sh, mx, my, display, - grid_step=20, line_color=(255, 60, 60, 235), - label_every=60, major_only_labels=True) - buf2 = io.BytesIO(); grid2.save(buf2, format="JPEG", quality=max(quality, 78), optimize=True) - data_grid2 = base64.b64encode(buf2.getvalue()).decode() - - return { - "data": data_clean, - "data_grid": data_grid1, - "data_grid2": data_grid2, - "data_grid3": data_grid2, # توافقية رجعية: نفس قيمة data_grid2 لأي كود قديم يقرأ data_grid3 - "width": ow, "height": oh, "mouse_x": mx, "mouse_y": my, - "is_placeholder": is_placeholder, # True إذا فشلت الطرق الثلاث وهذه صورة "Screenshot failed" - } - - -def _frame_hash(data: str) -> str: - """ - إصلاح: كان يُحسب على أول 2000 حرف فقط من بيانات الصورة، مما يعني أي تغيير - بصري يقع خارج الجزء الممثَّل ضمن تلك الأحرف الأولى (مثل تحديد نص شريط العنوان - بعد Ctrl+L، أو أي تغيير طفيف/في منطقة لا تتوافق مع بداية ترميز base64) لا - يُكتشف أبداً، فتُحجب الصورة الجديدة عبر delta suppression رغم اختلافها فعلياً - عن السابقة، ولا تصل أي لقطة محدَّثة للعميل. الحل: حساب MD5 على كامل البيانات - بدل تقطيعها — هذا سريع جداً (أقل من مللي ثانية حتى لصور كبيرة) فلا تكلفة أداء - حقيقية�� ويضمن اكتشاف أي تغيير بصري حقيقي في أي مكان من الصورة. - """ - return hashlib.md5(data.encode()).hexdigest() - - -# ════════════════════════════════════════════════════════════════ -# ── SafeSearch ─────────────────────────────────────────────── -# ════════════════════════════════════════════════════════════════ - -_BING_RE = re.compile(r"https?://(?:www\.)?bing\.com[^\s'\"]*") -_DDG_RE = re.compile(r"https?://(?:www\.)?duckduckgo\.com[^\s'\"]*") -_GOOG_RE = re.compile(r"https?://(?:www\.)?google\.[a-z.]+[^\s'\"]*") - -def _safe_search(text: str) -> str: - def _bing(m): - u = m.group(0) - return re.sub(r"adlt=\w+", "adlt=strict", u) if "adlt=" in u else u + ("&" if "?" in u else "?") + "adlt=strict" - def _ddg(m): - u = m.group(0) - return re.sub(r"kp=\d", "kp=1", u) if "kp=" in u else u + ("&" if "?" in u else "?") + "kp=1" - def _goog(m): - u = m.group(0) - return re.sub(r"safe=\w+", "safe=strict", u) if "safe=" in u else u + ("&" if "?" in u else "?") + "safe=strict" - return _GOOG_RE.sub(_goog, _DDG_RE.sub(_ddg, _BING_RE.sub(_bing, text))) - - -# ════════════════════════════════════════════════════════════════ -# ── Terminal Execution (per-display) ───────────────────────── -# ════════════════════════════════════════════════════════════════ - -_PKILL_RE = re.compile(r"\b(pkill|killall)\s+(-9\s+)?-f\s+(['\"]?)([a-zA-Z0-9_./-]+)\3") - -def _sanitize_pkill(cmd: str) -> str: - def _fix(m): - tool, d9, _q, p = m.group(1), m.group(2) or "", m.group(3), m.group(4) - sp = f"[{p[0]}]{p[1:]}" if len(p) > 1 else p - return f"{tool} {d9}-f '{sp}'" - return _PKILL_RE.sub(_fix, cmd) - - -async def run_cmd(cmd: str, display: str, timeout: int = 60) -> dict: - """ينفّذ أمر bash مع DISPLAY خاص بالمستخدم.""" - cmd = _sanitize_pkill(_safe_search(cmd)) - async with _terminal_sem: - env = {**os.environ, "DISPLAY": display, - "PYTHONIOENCODING": "utf-8", "LANG": "en_US.UTF-8"} - - def _exec(): - try: - r = subprocess.run( - cmd, shell=True, capture_output=True, - text=True, timeout=timeout, env=env, executable="/bin/bash" - ) - return {"stdout": r.stdout[-15000:], "stderr": r.stderr[-3000:], "returncode": r.returncode} - except subprocess.TimeoutExpired: - return {"stdout": "", "stderr": f"⏱️ Timeout {timeout}s", "returncode": -1} - except Exception as e: - return {"stdout": "", "stderr": str(e), "returncode": -1} - - return await asyncio.to_thread(_exec) - - -# ════════════════════════════════════════════════════════════════ -# ── xdotool helpers (per-display) ──────────────────────────── -# ════════════════════════════════════════════════════════════════ - -async def xdo(args: list, display: str, timeout: int = 10) -> dict: - env = {**os.environ, "DISPLAY": display} - def _run(): - r = subprocess.run( - ["xdotool"] + args, env=env, - timeout=timeout, capture_output=True, text=True - ) - return {"rc": r.returncode, "out": r.stdout, "err": r.stderr} - return await asyncio.to_thread(_run) - - -async def type_smart(text: str, display: str) -> dict: - """كتابة نص ذكية بـ display محدد.""" - has_arabic = bool(re.search(r'[\u0600-\u06FF]', text)) - env = {**os.environ, "DISPLAY": display} - - if has_arabic: - def _paste(): - p = subprocess.Popen( - ["xclip", "-selection", "clipboard"], - stdin=subprocess.PIPE, env=env - ) - p.communicate(text.encode("utf-8")) - await asyncio.to_thread(_paste) - await asyncio.sleep(0.15) - await xdo(["key", "--clearmodifiers", "ctrl+v"], display) - return {"method": "clipboard+paste"} - else: - r = await xdo(["type", "--clearmodifiers", "--delay", "25", text], display) - return {"method": "xdotool", "rc": r["rc"]} - - -# ════════════════════════════════════════════════════════════════ -# ── Search Sources ──────────────────────────────────────────── -# ════════════════════════════════════════════════════════════════ - -def _search_sources(query: str) -> list: - q = urllib.parse.quote_plus(query) - return [ - {"name": "DuckDuckGo Instant", - "cmd": f"curl -s --max-time 15 'https://api.duckduckgo.com/?q={q}&format=json&no_html=1&skip_disambig=1' | python3 -c \"import sys,json;d=json.load(sys.stdin);a=d.get('AbstractText','');r=d.get('RelatedTopics',[]);print('ANS:',a or 'none');[print('-',x.get('Text','')[:200]) for x in r[:6] if isinstance(x,dict)]\""}, - {"name": "Google News RSS", - "cmd": f"curl -sL --max-time 15 'https://news.google.com/rss/search?q={q}&hl=ar&gl=AR&ceid=AR:ar' | python3 -c \"import sys,re;x=sys.stdin.read();t=re.findall(r'
✅ Server RUNNING — Per-User Isolated Displays
-🌐 Browser: {BROWSER}
-👥 Active users: {n}
-📋 Display range: :{DISPLAY_BASE} → :{DISPLAY_MAX}
-| User ID | Display | Status |
|---|---|---|
| No active sessions | ||
Endpoints: /health · /ws?user_id=email (WebSocket)
-""" - - -# ════════════════════════════════════════════════════════════════ -# ── Action Handler (per-session) ───────────────────────────── -# ════════════════════════════════════════════════════════════════ - -def _format_elements_for_llm(elements: list) -> str: - """ - يُحوّل قائمة عناصر DOM لنص منظم يمكن للـ LLM قراءته واستخدامه - لاختيار العنصر المناسب بـ id — مثل ما يرجعه OmniParser/browser_view في Manus. - - مثال على الخرج: - [0] BUTTON "Search" at (640, 50) — 120×40px - [1] INPUT[text] "" at (400, 50) placeholder="Search..." - [2] LINK "Home" → /home at (100, 30) - [3] BUTTON "Login" at (1100, 30) — 80×35px - """ - if not elements: - return "No interactive elements found on this page." - - lines = [f"📋 Found {len(elements)} interactive elements:\n"] - - # أولاً: عناصر داخل viewport (مرئية) - in_view = [e for e in elements if e.get("in_viewport", True)] - out_view = [e for e in elements if not e.get("in_viewport", True)] - - if in_view: - lines.append("── Visible elements (in viewport) ──") - for el in in_view[:50]: # حد أقصى 50 لتجنب overflow - lines.append(_format_single_element(el)) - - if out_view: - lines.append(f"\n── Off-screen elements ({len(out_view)} total, scroll to see) ──") - for el in out_view[:20]: - lines.append(_format_single_element(el)) - - lines.append(f"\n💡 To click element N: use action='browser_click_element' with id=N") - lines.append(f"💡 To type into an input N: use action='browser_type_into_element' with id=N, text='...'") - lines.append(f"💡 To find an element by description: use action='browser_find_element' with query='...'") - return "\n".join(lines) - - -def _format_single_element(el: dict) -> str: - """يُنسّق عنصر واحد لخط واحد مقروء.""" - eid = el.get("id", "?") - tag = el.get("tag", "?").upper() - text = el.get("text", "").strip()[:60] - x, y = el.get("x", 0), el.get("y", 0) - w, h = el.get("w", 0), el.get("h", 0) - el_type = el.get("type", "") - href = el.get("href", "") - role = el.get("role", "") - - # تصنيف نوع العنصر - if tag == "A" and href: - kind = f"LINK → {href[:40]}" - elif tag == "INPUT": - kind = f"INPUT[{el_type or 'text'}]" - elif tag == "BUTTON" or role == "button": - kind = "BUTTON" - elif tag == "SELECT": - kind = "SELECT (dropdown)" - elif tag == "TEXTAREA": - kind = "TEXTAREA" - else: - kind = role.upper() if role and role != tag.lower() else tag - - text_part = f'"{text}"' if text else '(no text)' - coord_part = f"at ({x},{y}) {w}×{h}px" - - return f" [{eid}] {kind} {text_part} {coord_part}" - - -async def handle_action(ws: WebSocket, msg: dict, sess: dict): - action = msg.get("action", "") - data = msg.get("data", {}) - display = sess["display"] - - async def send(obj): - try: - await ws.send_text(json.dumps(obj, ensure_ascii=False)) - except Exception: - pass - - # ── إصلاح جذري: قفل خاص بكل session لمنع تراكم/تضارب عمليات capture المتزامنة ── - # المشكلة الأصلية: action == "screenshot" كان يُنفَّذ بـ await مباشر ضمن حلقة - # while True الرئيسية في websocket_endpoint، فإذا تأخرت capture_with_grid (قد تصل - # لدقيقة ونصف في أسوأ سيناريو فشل/تأخر scrot+import+ffmpeg)، تبقى الحلقة كلها - # محجوبة ولا تستقبل أي رسالة عميل جديدة (نقرة، طلب screenshot آخر، إلخ) حتى تنتهي. - # الحل: كل طلب screenshot يُشغَّل فوراً عبر create_task (لا يحجب الحلقة أبداً). - # - # ── إصلاح ثانٍ (مهم جداً): طلب screenshot الصريح من العميل له أولوية مطلقة ── - # المشكلة المكتشفة لاحقاً: لو استُخدم قفل واحد مشترك بين shot_bg (تلقائية، بعد - # كل فعل كنقرة/فتح تبويب) و shot_explicit (صريحة، يطلبها العميل مباشرة بعد كل - # خطوة لتحديث الصورة المعروضة)، فإن shot_explicit قد تنتظر خلف shot_bg طويلاً - # إذا كانت الأخيرة قد بدأت فعلاً وعلقت داخل capture_with_grid (محاولات فاشلة - # متتالية). هذا يجعل العميل يرى "لم تصل لقطة شاشة" بشكل متكرر حتى بعد نجاح - # الخطوة الفعلية (مثل open_tab)، لأن طلبه العاجل كان يصطف خلف عملية تلقائية بطيئة. - # الحل: shot_explicit (الصريحة فقط) تحاول الحصول على القفل لفترة قصيرة جداً - # (1 ثانية)، وإن لم تنجح (لأن shot_bg تستخدمه)، تُنفَّذ التصوير مباشرة بدون قفل - # بدل الانتظار — لأن استجابة العميل الفورية أهم من تفادي تزاحم CPU عرضي بسيط، - # وقراءة الشاشة (X11) عملية قراءة فقط لا تُسبب أي تلف بيانات عند التزاحم. - if "_shot_lock" not in sess: - sess["_shot_lock"] = asyncio.Lock() - _shot_lock = sess["_shot_lock"] - - async def _safe_capture(scale, quality, force_mx=None, force_my=None): - """ - تُستخدم من shot_bg (التلقائية) — تنتظر القفل ثم تنتظر التقاط الشاشة - حتى ينتهي فعلياً، بدون أي سقف زمني يقطعها. هذا يضمن أن العملية لن - تُقطع أبداً في منتصفها وتُرجع فراغاً؛ ستكمل حتى تنجح (أو تفشل كل - المحاولات الداخلية في capture_with_grid وتُرجع placeholder صريح). - """ - async with _shot_lock: - return await asyncio.to_thread( - capture_with_grid, display, scale, quality, force_mx, force_my - ) - - async def _priority_capture(scale, quality, force_mx=None, force_my=None): - """ - تُستخدم من shot_explicit (الطلب الصريح من العميل) فقط — أولوية قصوى. - تحاول الحصول على القفل لمدة قصيرة (1 ثانية) فقط؛ إن لم تنجح (القفل - محجوز من shot_bg تلقائية بطيئة)، تُنفَّذ التصوير فوراً بدون قفل بدل - الانتظار خلف عملية أخرى — العميل يجب أن يحصل على رد سريع دائماً. - بعد الحصول على القفل (أو تجاوزه)، لا يوجد أي سقف زمني على عملية - الالتقاط نفسها — تنتظر حتى تكتمل فعلياً بدل أن تُقطع في المنتصف. - """ - try: - await asyncio.wait_for(_shot_lock.acquire(), timeout=1.0) - try: - return await asyncio.to_thread( - capture_with_grid, display, scale, quality, force_mx, force_my - ) - finally: - _shot_lock.release() - except asyncio.TimeoutError: - # لم نحصل على القفل بسرعة كافية — ننفّذ التصوير مباشرة بدون قفل - # بدل الانتظار خلف عملية أخرى، لكن بدون أي سقف زمني على الالتقاط - return await asyncio.to_thread( - capture_with_grid, display, scale, quality, force_mx, force_my - ) - - async def shot_bg(label: str = "", delay: float = 0.5, - force_mx: int | None = None, force_my: int | None = None, - extra_shot_delay: float | None = None): - """ - auto_shot في الخلفية — مع rate limit + delta suppression. - إذا تم تمرير extra_shot_delay، تُرسل لقطة ثانية إضافية بعد ذلك التأخير - الإضافي (محسوباً من وقت انتهاء اللقطة الأولى) — هذا يغطي حالة النقر على - رابط ينقل لصفحة جديدة قد لا تكتمل تحميلها خلال التأخير الأولي القصير - (مثلاً صفحة بطيئة على سيرفر محدود الموارد)، دون الحاجة لمعرفة مسبقة بنوع - العنصر المنقور عليه. تتجاهل delta suppression لهذه اللقطة الثانية تحديداً - لأن المحتوى متوقع أن يكون مختلفاً (صفحة تحمّلت أكثر) حتى لو تشابه الـ hash - جزئياً مع خلفية مشابهة. - """ - try: - await asyncio.sleep(delay) - now = time.time() - if now - sess.get("last_bg_shot_ts", 0) < 0.3: - await asyncio.sleep(0.3 - (now - sess["last_bg_shot_ts"])) - - result = await _safe_capture(0.65, 72, force_mx, force_my) - sess["last_bg_shot_ts"] = time.time() - # ── مساهمة في عدّاد الفشل المشترك مع shot_explicit ── - # shot_bg تبقى صامتة عند الفشل (لا ترسل شيئاً — هذا سلوكها الأصلي - # المقصود، فهي مجرد محاولة إضافية سريعة وليست خط الدفاع الموثوق)، - # لكنها الآن تُحدّث نفس عدّاد الفشل المتتالي الذي يستخدمه - # shot_explicit. هذا يعني: لو الـ display معطوب فعلاً، العدّاد - # يبدأ بالارتفاع من أول لقطة تلقائية فاشلة، فتصل shot_explicit - # (التي دائماً تُطلَب صراحة بعدها من العميل) للعتبة وتُفعّل الـ - # reset التلقائي أسرع، بدل الانتظار حتى تتراكم فشلتان داخل - # shot_explicit نفسها فقط. - if result.get("data") and not result.get("is_placeholder"): - sess["consecutive_capture_failures"] = 0 - fh = _frame_hash(result["data"]) - if fh != sess.get("last_bg_hash", ""): - sess["last_bg_hash"] = fh - await send({ - "type": "screenshot", - "data": result["data"], - "data_grid": result.get("data_grid", ""), - "data_grid2": result.get("data_grid2", ""), - "data_grid3": result.get("data_grid3", ""), - "ts": int(time.time() * 1000), - "auto": True, "label": label, - "screen_width": result["width"], - "screen_height": result["height"], - "mouse_x": result["mouse_x"], - "mouse_y": result["mouse_y"], - "has_grid": True, - }) - else: - sess["consecutive_capture_failures"] = sess.get("consecutive_capture_failures", 0) + 1 - - if extra_shot_delay: - await asyncio.sleep(extra_shot_delay) - result2 = await _safe_capture(0.65, 72, force_mx, force_my) - sess["last_bg_shot_ts"] = time.time() - if result2["data"]: - fh2 = _frame_hash(result2["data"]) - sess["last_bg_hash"] = fh2 - await send({ - "type": "screenshot", - "data": result2["data"], - "data_grid": result2.get("data_grid", ""), - "data_grid2": result2.get("data_grid2", ""), - "data_grid3": result2.get("data_grid3", ""), - "ts": int(time.time() * 1000), - "auto": True, "label": f"{label} (delayed)", - "screen_width": result2["width"], - "screen_height": result2["height"], - "mouse_x": result2["mouse_x"], - "mouse_y": result2["mouse_y"], - "has_grid": True, - }) - except Exception as e: - # لقطة خلفية فاشلة لا يجب أن توقف الجلسة أو تظهر بصمت في اللوق فقط. - print(f"[shot_bg:{display}] ⚠️ {e}") - - # ── عتبة الفشل المتتالي قبل تفعيل reset تلقائي كامل للـ display ── - # فشلة واحدة عابرة (ضغط CPU لحظي، تزاحم مع shot_bg) لا تستدعي reset — - # لكن فشلتان متتاليتان فعليتان (كلتاهما رجعتا is_placeholder=True) تعني - # أن هناك مشكلة حقيقية في الـ display نفسه (Xvfb ميت/عالق)، وليس عطلاً - # عابراً، فنعيد ضبط الكمبيوتر تلقائياً بدل الاستسلام وإرسال "error". - _AUTO_RESET_FAILURE_THRESHOLD = 2 - # حد أقصى لعدد جولات (محاولة التقاط → عند الفشل المتكرر → reset → إعادة - # محاولة) داخل نفس طلب screenshot الواحد، لضمان أن الحلقة تنتهي دائماً - # ولو استمر الفشل حتى بعد عدة عمليات reset (بيئة معطوبة جذرياً)، بدل أن - # تدور للأبد بصمت. هذا ليس "timeout" زمنياً — هو س��ف على عدد المحاولات. - _MAX_CAPTURE_ROUNDS = 6 - - async def shot_explicit(label: str = ""): - """screenshot صريح — يُرسل دائماً بدون delta suppression، وله أولوية مطلقة - على أي عملية shot_bg تلقائية جارية (لا ينتظر خلفها طويلاً). - - ── إصلاح جذري v14: لا "error" نهائي بعد الآن — بل reset تلقائي كامل ── - سابقاً: فشل التقاط واحد (نتيجة فارغة) كان يُرسل "error" فوراً للعميل، - فيبقى العميل يعيد طلب screenshot بلا نهاية بينما السيرفر لا يفعل شيئاً - مختلفاً في كل مرة (نفس الـ display المعطوب، بلا أي محاولة إصلاح فعلية). - - الآن: فشل واحد → إعادة محاولة فورية (قد يكون عابراً). فشلان متتاليان - فعليان (كلاهما is_placeholder=True) → reset تلقائي كامل لهذا المستخدم - (نفس ما يفعله المستخدم يدوياً بزر "إعادة الضبط") ثم إعادة المحاولة - تلقائياً على الـ display الجديد النظيف — كل هذا داخل نفس الطلب، بلا - أي تدخل من العميل. رسالة "error" الصريحة لم تعد تُرسَل أبداً من هذا - المسار؛ العميل يستقبل إما screenshot ناجحاً حقيقياً، أو (في أسوأ حالة - نادرة جداً: بيئة معطوبة جذرياً حتى بعد عدة محاولات reset) رسالة - "computer_reset" تخبره أن الجلسة بدأت من جديد ليطلب screenshot تالياً - بنفسه — أبداً حالة صمت لا نهائية. - """ - for _round in range(_MAX_CAPTURE_ROUNDS): - try: - result = await _priority_capture(0.65, 75) - got_real_frame = bool(result and result.get("data") and not result.get("is_placeholder")) - - if got_real_frame: - sess["consecutive_capture_failures"] = 0 - sess["last_bg_hash"] = _frame_hash(result["data"]) - await send({ - "type": "screenshot", - "data": result["data"], - "data_grid": result.get("data_grid", ""), - "data_grid2": result.get("data_grid2", ""), - "data_grid3": result.get("data_grid3", ""), - "ts": int(time.time() * 1000), - "auto": False, "label": label, - "screen_width": result["width"], - "screen_height": result["height"], - "mouse_x": result["mouse_x"], - "mouse_y": result["mouse_y"], - "has_grid": True, - }) - return - - # ── فشل هذه الجولة (نتيجة فارغة تماماً، أو placeholder) ── - sess["consecutive_capture_failures"] = sess.get("consecutive_capture_failures", 0) + 1 - fail_count = sess["consecutive_capture_failures"] - print(f"[shot_explicit:{display}] ⚠️ capture failed (round {_round+1}/{_MAX_CAPTURE_ROUNDS}, " - f"consecutive={fail_count})") - - if fail_count >= _AUTO_RESET_FAILURE_THRESHOLD: - # ── مشكلة حقيقية وليست عابرة — reset تلقائي كامل ── - print(f"[shot_explicit:{display}] 🔄 auto-reset triggered after {fail_count} consecutive failures") - await send({ - "type": "computer_reset", - "msg": "🔄 إعادة ضبط تلقائية بعد فشل التقاط متكرر — الشاشة تبدأ من جديد نظيفة", - "auto": True, - }) - await reset_user_computer(sess["user_id"]) - # بعد الـ reset الكامل، أمهل Xvfb الجديد لحظة قصيرة ليستقر - # قبل أول محاولة التقاط عليه (نفس المهلة المستخدمة في مسار - # reset_computer اليدوي أعلاه لضمان اتساق السلوك). - await asyncio.sleep(2.0) - # الجولة التالية من هذه الحلقة نفسها ستعيد المحاولة على - # الـ display الجديد النظيف تلقائياً — بلا حاجة لأي طلب - # إضافي من العميل ولا لإرسال "error" له في الأثناء. - continue - - # فشل عابر فقط (أول فشلة) — أعد المحاولة فوراً بلا reset - await asyncio.sleep(0.4) - continue - - except Exception as e: - # استثناء غير متوقع تماماً (وليس مجرد نتيجة فارغة) — يُعامَل - # بنفس منطق الفشل أعلاه بدل الاستسلام الفوري، لأن أسباباً - # كهذه (اتصال X عالق، انهيار عملية) هي بالضبط ما يُصلحه reset. - sess["consecutive_capture_failures"] = sess.get("consecutive_capture_failures", 0) + 1 - fail_count = sess["consecutive_capture_failures"] - print(f"[shot_explicit:{display}] ❌ EXCEPTION (round {_round+1}/{_MAX_CAPTURE_ROUNDS}, " - f"consecutive={fail_count}): {e}") - if fail_count >= _AUTO_RESET_FAILURE_THRESHOLD: - try: - await send({ - "type": "computer_reset", - "msg": "🔄 إعادة ضبط تلقائية بعد خطأ متكرر في الالتقاط — الشاشة تبدأ من جديد نظيفة", - "auto": True, - }) - except Exception: - pass - try: - await reset_user_computer(sess["user_id"]) - await asyncio.sleep(2.0) - except Exception as reset_ex: - print(f"[shot_explicit:{display}] ❌ reset itself failed: {reset_ex}") - continue - - # ── وصلنا لحد _MAX_CAPTURE_ROUNDS جولة دون أي نجاح حقيقي، حتى بعد - # محاولات reset متعددة — بيئة معطوبة جذرياً (نادر جداً). بدل الصمت - # التام، نُعلم العميل صراحة أن الجلسة أُعيد ضبطها ليطلب screenshot - # تالياً بنفسه بدل انتظار frame لن يصل أبداً في هذا الطلب تحديداً. - print(f"[shot_explicit:{display}] ❌ giving up after {_MAX_CAPTURE_ROUNDS} rounds — " - f"environment appears broken even after reset attempts") - try: - await send({ - "type": "computer_reset", - "msg": "⚠️ استمرت مشكلة الالتقاط رغم إعادة الضبط — الرجاء طلب لقطة شاشة جديدة", - "auto": True, - }) - except Exception: - pass - - # ── reset_computer: إعادة ضبط الكمبيوتر كأنه جديد ───────── - if action == "reset_computer": - user_id = sess["user_id"] - await send({"type": "ack", "action": "reset_computer", "status": "resetting"}) - await reset_user_computer(user_id) - # أرسل لقطة شاشة بعد الريستارت - await asyncio.sleep(2.0) - await shot_explicit("after reset") - await send({"type": "computer_reset", "msg": "✅ تم إعادة ضبط الكمبيوتر — الشاشة جديدة تماماً"}) - return - - # ── screenshot ──────────────────────────────────────────── - # إصلاح: create_task بدل await مباشر — لا يحجب حلقة استقبال الرسائل أبداً، - # حتى لو تأخرت capture_with_grid لأي سبب (CPU ضعيف، X11 بطيء، إلخ). - # هذا يضمن أن أي إجراء لاحق (نقرة، طباعة، طلب screenshot آخر) يصل ويُعالَج - # فوراً دون انتظار اكتمال هذه اللقطة أولاً. - if action == "screenshot": - # ── طبقة حماية إضافية: حتى لو حصل خطأ غير متوقع تماماً داخل - # shot_explicit تفلت من الـ try/except الداخلي (مثل CancelledError - # أو خطأ عند إنشاء الـ task نفسه)، هذا الغلاف الخارجي يضمن تسجيله - # في اللوق بدل أن يختفي بصمت كـ "unhandled task exception" في event - # loop بايثون (وهذا كان له سلوك افتراضي بصامت في بعض إعدادات uvicorn). - _task = asyncio.create_task(shot_explicit("explicit screenshot")) - def _on_shot_done(t: asyncio.Task): - exc = t.exception() if not t.cancelled() else None - if exc: - print(f"[screenshot-task:{display}] ❌ unhandled exception: {exc}") - _task.add_done_callback(_on_shot_done) - - # ── terminal ───────────────────────────────────────────── - elif action == "terminal": - cmd = data.get("cmd", "") - if not cmd: - await send({"type": "terminal_result", "stdout": "", "stderr": "no cmd", "returncode": -1}) - return - res = await run_cmd_smart(cmd, display, int(data.get("timeout", 60))) - await send({ - "type": "terminal_result", - "cmd": cmd, - "stdout": res["stdout"], - "stderr": res.get("stderr", ""), - "returncode": res["returncode"], - }) - # ── تعطيل الالتقاط التلقائي: لا تُؤخذ لقطة شاشة إلا بطلب صريح ── - - # ── mouse_move ─────────────────────────────────────────── - elif action == "mouse_move": - x, y = int(data.get("x", 0)), int(data.get("y", 0)) - await xdo(["mousemove", "--sync", str(x), str(y)], display) - await send({"type": "ack", "action": "mouse_move", "x": x, "y": y}) - # ── تعطيل الالتقاط التلقائي ── - - # ── mouse_click ────────────────────────────────────────── - elif action == "mouse_click": - x, y = int(data.get("x", 0)), int(data.get("y", 0)) - btn = {"left": "1", "middle": "2", "right": "3"}.get(data.get("button", "left"), "1") - double = data.get("double", False) - await xdo(["mousemove", "--sync", str(x), str(y)], display) - await asyncio.sleep(0.07) - if double: - await xdo(["click", "--repeat", "2", "--delay", "100", btn], display) - else: - await xdo(["click", btn], display) - btn_name = {"1": "left", "2": "middle", "3": "right"}.get(btn, "left") - await send({"type": "ack", "action": "mouse_click", "x": x, "y": y, "button": btn_name}) - # ── تعطيل الالتقاط التلقائي ── - - # ── mouse_drag ─────────────────────────────────────────── - elif action == "mouse_drag": - x1, y1 = int(data.get("x1", 0)), int(data.get("y1", 0)) - x2, y2 = int(data.get("x2", 0)), int(data.get("y2", 0)) - await xdo(["mousemove", str(x1), str(y1)], display) - await xdo(["mousedown", "1"], display) - await asyncio.sleep(0.1) - await xdo(["mousemove", str(x2), str(y2)], display) - await asyncio.sleep(0.1) - await xdo(["mouseup", "1"], display) - await send({"type": "ack", "action": "mouse_drag"}) - # ── تعطيل الالتقاط التلقائي ── - - # ── keyboard_type ──────────────────────────────────────── - elif action == "keyboard_type": - text = _safe_search(data.get("text", "")) - if text: - res = await type_smart(text, display) - await send({"type": "ack", "action": "keyboard_type", "method": res["method"]}) - # ── تعطيل الالتقاط التلقائي ── - - # ── keyboard_hotkey ────────────────────────────────────── - elif action == "keyboard_hotkey": - keys = data.get("keys", []) - if keys: - await xdo(["key", "--clearmodifiers", "+".join(keys)], display) - await send({"type": "ack", "action": "keyboard_hotkey", "keys": keys}) - # ── تعطيل الالتقاط التلقائي ── - - # ── keyboard_press ─────────────────────────────────────── - elif action == "keyboard_press": - key = data.get("key", "") - if key: - await xdo(["key", "--clearmodifiers", key], display) - await send({"type": "ack", "action": "keyboard_press"}) - # ── تعطيل الالتقاط التلقائي ── - - # ── scroll ─────────────────────────────────────────────── - elif action == "scroll": - x, y = int(data.get("x", 960)), int(data.get("y", 540)) - clicks = max(-5, min(5, int(data.get("clicks", 3)))) - btn = "4" if clicks > 0 else "5" - await xdo(["mousemove", str(x), str(y)], display) - for _ in range(abs(clicks)): - await xdo(["click", btn], display) - await asyncio.sleep(0.025) - await send({"type": "ack", "action": "scroll", "clicks": clicks}) - # ── تعطيل الالتقاط التلقائي ── - - # ── clipboard_write ────────────────────────────────────── - elif action == "clipboard_write": - text = data.get("text", "") - env = {**os.environ, "DISPLAY": display} - def _clip(): - p = subprocess.Popen(["xclip", "-selection", "clipboard"], stdin=subprocess.PIPE, env=env) - p.communicate(text.encode("utf-8")) - await asyncio.to_thread(_clip) - await send({"type": "ack", "action": "clipboard_write", "length": len(text)}) - - # ── clipboard_read ─────────────────────────────────────── - elif action == "clipboard_read": - res = await run_cmd("xclip -selection clipboard -o", display, 5) - await send({"type": "clipboard_content", "text": res["stdout"]}) - - # ── paste ──────────────────────────────────────────────── - elif action == "paste": - text = data.get("text", "") - if text: - env = {**os.environ, "DISPLAY": display} - def _clip2(): - p = subprocess.Popen(["xclip", "-selection", "clipboard"], stdin=subprocess.PIPE, env=env) - p.communicate(text.encode("utf-8")) - await asyncio.to_thread(_clip2) - await asyncio.sleep(0.1) - await xdo(["key", "--clearmodifiers", "ctrl+v"], display) - await send({"type": "ack", "action": "paste"}) - # ── تعطيل الالتقاط التلقائي ── - - # ── open_app ───────────────────────────────────────────── - elif action == "open_app": - cmd = _safe_search(data.get("cmd", "")) - if not cmd: - await send({"type": "ack", "action": "open_app"}) - return - if "firefox" in cmd.lower(): - profile_dir = _ensure_firefox_profile(sess["user_id"], display) - # قتل فوري وقوي (SIGKILL) بدل pkill اللطيف: يمنع فايرفوكس من - # "التقاط" حالة نصف-مفتوحة تُسجَّل لاحقاً كتحطّم (crash) يُنتج - # شاشة "Restore Session" في المرة القادمة. - await run_cmd( - "pkill -9 -f '[f]irefox' 2>/dev/null; sleep 0.6; echo CLEANED", - display, timeout=8 - ) - _wipe_firefox_session_data(profile_dir) - # حقن --profile داخل أمر فايرفوكس (إن لم يكن محقوناً مسبقاً) - if "--profile" not in cmd and "-P " not in cmd: - flags_str = " ".join(FIREFOX_CLI_FLAGS) - cmd = cmd.replace("firefox", f"firefox --profile '{profile_dir}' {flags_str}", 1) - env = {**os.environ, "DISPLAY": display, **(FIREFOX_ENV_EXTRA if "firefox" in cmd.lower() else {})} - proc = subprocess.Popen(cmd, shell=True, env=env, - stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) - if any(b in cmd for b in ("firefox", "chromium", "chrome")): - with _display_lock: - if sess["user_id"] in _user_displays: - _user_displays[sess["user_id"]]["browser_proc"] = proc - await send({"type": "ack", "action": "open_app", "cmd": cmd}) - # ── تعطيل الالتقاط التلقائي بعد فتح التطبيق ── - - # ── open_browser ───────────────────────────────────────── - elif action == "open_browser": - url = _safe_search(data.get("url", "") or "about:blank") - - if BROWSER_SUPPORTS_CDP: - # v12: Chromium مع CDP — يُلغي حاجة Firefox profiles بالكامل - port = _get_cdp_port(display) - await run_cmd( - f"pkill -9 -f '[r]emote-debugging-port={port}' 2>/dev/null; " - f"pkill -9 -f '[c]hromium' 2>/dev/null; " - f"pkill -9 -f '[g]oogle-chrome' 2>/dev/null; sleep 0.5; echo CLEANED", - display, timeout=8 - ) - # امسح CDPSession القديمة - with _cdp_sessions_lock: - old_cdp = _cdp_sessions.pop(display, None) - if old_cdp: - asyncio.create_task(old_cdp.close()) - - proc = await asyncio.to_thread(_start_chromium_with_cdp, display, url) - with _display_lock: - if sess["user_id"] in _user_displays: - _user_displays[sess["user_id"]]["browser_proc"] = proc - _user_displays[sess["user_id"]]["last_elements"] = [] - await send({"type": "ack", "action": "open_browser", "url": url, - "cdp": True, "cdp_port": port}) - else: - # fallback: Firefox (سلوك v11 الأصلي) - profile_dir = _ensure_firefox_profile(sess["user_id"], display) - await run_cmd( - "pkill -9 -f '[f]irefox' 2>/dev/null; sleep 0.6; echo CLEANED", - display, timeout=8 - ) - _wipe_firefox_session_data(profile_dir) - env_ff = {**os.environ, "DISPLAY": display, **FIREFOX_ENV_EXTRA} - proc = subprocess.Popen( - [BROWSER, "--profile", profile_dir, *FIREFOX_CLI_FLAGS, url], - env=env_ff, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL - ) - with _display_lock: - if sess["user_id"] in _user_displays: - _user_displays[sess["user_id"]]["browser_proc"] = proc - await send({"type": "ack", "action": "open_browser", "url": url, "cdp": False}) - # ── تعطيل الالتقاط التلقائي بعد فتح المتصفح ── - - # ── open_tab ───────────────────────────────────────────── - elif action == "open_tab": - # إصلاح: زيادة التأخيرات بين كل خطوة فرعية هنا (Ctrl+T → Ctrl+L → كتابة الرابط) - # المشكلة الأصلية: التأخير القصير (0.4s) بعد Ctrl+T لم يكن كافياً دائماً، خصوصاً - # على سيرفر محدود الموارد قد يستغرق فايرفوكس فيه وقتاً أطول لإنشاء التبويب الجديد - # فعلياً وتحويل التركيز (focus) إليه. إذا وصل Ctrl+L والكتابة قبل اكتمال ذلك، - # فإنها قد تذهب للتبويب القديم بالخطأ، فلا تنتقل الصفحة فعلياً رغم نجاح الأمر - # ظاهرياً (بدون أي رسالة خطأ)، فيبقى المستخدم يرى نفس الصفحة القديمة للأبد. - url = _safe_search(data.get("url", "") or "about:blank") - await xdo(["key", "--clearmodifiers", "ctrl+t"], display) - await asyncio.sleep(0.8) - await xdo(["key", "--clearmodifiers", "ctrl+l"], display) - await asyncio.sleep(0.3) - await type_smart(url, display) - await asyncio.sleep(0.3) - await xdo(["key", "--clearmodifiers", "Return"], display) - await send({"type": "ack", "action": "open_tab", "url": url}) - # ── تعطيل الالتقاط التلقائي ── - - # ── close_tab ──────────────────────────────────────────── - elif action == "close_tab": - await xdo(["key", "--clearmodifiers", "ctrl+w"], display) - await send({"type": "ack", "action": "close_tab"}) - # ── تعطيل الالتقاط التلقائي ── - - # ── browser_back ───────────────────────────────────────── - elif action == "browser_back": - await xdo(["key", "--clearmodifiers", "alt+Left"], display) - await send({"type": "ack", "action": "browser_back"}) - # ── تعطيل الالتقاط التلقائي ── - - # ── browser_forward ────────────────────────────────────── - elif action == "browser_forward": - await xdo(["key", "--clearmodifiers", "alt+Right"], display) - await send({"type": "ack", "action": "browser_forward"}) - # ── تعطيل الالتقاط التلقائي ── - - # ── browser_search ─────────────────────────────────────── - elif action == "browser_search": - url = _safe_search(data.get("url", "") or data.get("query", "")) - await xdo(["key", "--clearmodifiers", "ctrl+l"], display) - await asyncio.sleep(0.2) - await type_smart(url, display) - await asyncio.sleep(0.15) - await xdo(["key", "--clearmodifiers", "Return"], display) - await send({"type": "ack", "action": "browser_search"}) - # ── تعطيل الالتقاط التلقائي ── - - # ── screen_info ────────────────────────────────────────── - elif action == "screen_info": - w, h = await asyncio.to_thread(_get_screen_size, display) - mx, my = await asyncio.to_thread(_get_mouse_pos, display) - await send({ - "type": "screen_info", - "width": w, "height": h, - "mouse_x": mx, "mouse_y": my, - "browser": BROWSER, - "display": display, - }) - - # ══════════════════════════════════════════════════════════ - # ── v12 CDP ACTIONS — مستوحاة من Manus browser_* tools ── - # ══════════════════════════════════════════════════════════ - - # ── browser_get_elements ───────────────────────────────── - # مثل OmniParser لكن بدون GPU: يجمع كل العناصر القابلة للنقر - # من DOM مباشرة مع إحداثياتها الدقيقة 100% - elif action == "browser_get_elements": - if not BROWSER_SUPPORTS_CDP: - await send({ - "type": "browser_elements", - "error": "CDP requires Chromium (Firefox not supported)", - "elements": [], - }) - return - - result = await cdp_get_elements(display) - elements = result.get("elements", []) - - # خزّن قائمة العناصر + رابط الصفحة وقت الالتقاط، للاستخدام في - # click/type/scroll لاحقاً مع التحقق من عدم تغيّر الصفحة (staleness guard) - with _display_lock: - if sess["user_id"] in _user_displays: - _user_displays[sess["user_id"]]["last_elements"] = elements - _user_displays[sess["user_id"]]["last_elements_url"] = result.get("page_url", "") - - await send({ - "type": "browser_elements", - "elements": elements, - "page_title": result.get("page_title", ""), - "page_url": result.get("page_url", ""), - "viewport_w": result.get("viewport_w", 1280), - "viewport_h": result.get("viewport_h", 800), - "scroll_y": result.get("scroll_y", 0), - "total_height": result.get("total_height", 0), - "count": len(elements), - "error": result.get("error"), - # ملخص نصي للـ LLM (مثل ما يرجع Manus في browser_view) - "summary": _format_elements_for_llm(elements), - }) - - # ── browser_click_element ──────────────────────────────── - # نقر بـ index مباشرة — مثل Manus: browser_click(index=5) - # الإحداثيات مأخوذة من آخر استدعاء لـ browser_get_elements - elif action == "browser_click_element": - element_id = int(data.get("id", -1)) - - with _display_lock: - user_disp = _user_displays.get(sess["user_id"], {}) - elements = user_disp.get("last_elements", []) - cached_url = user_disp.get("last_elements_url", "") - - if element_id < 0 or not elements: - await send({ - "type": "error", - "msg": "Call browser_get_elements first, then use the returned element id", - }) - return - - if await _elements_are_stale(display, cached_url): - await send({ - "type": "error", - "msg": "Page changed since last browser_get_elements — element coordinates are stale. Call browser_get_elements again before clicking.", - "stale": True, - }) - return - - click_info = await cdp_click_element(display, element_id, elements) - if "error" in click_info: - await send({"type": "error", "msg": click_info["error"]}) - return - - # نفّذ النقر بالإحداثيات الحقيقية عبر xdotool - x, y = click_info["x"], click_info["y"] - await xdo(["mousemove", "--sync", str(x), str(y)], display) - await asyncio.sleep(0.07) - await xdo(["click", "1"], display) - - await send({ - "type": "ack", - "action": "browser_click_element", - "element_id": element_id, - "x": x, "y": y, - "tag": click_info.get("tag", ""), - "text": click_info.get("text", ""), - }) - asyncio.create_task(shot_bg("after element click", delay=0.8, extra_shot_delay=2.0)) - - # ── browser_type_into_element ──────────────────────────── - # يكتب مباشرة في حقل بـ index — يدمج click + keyboard_type في نداء - # واحد (مثل Manus browser_input(index, text)) بدل خطوتين منفصلتين - elif action == "browser_type_into_element": - element_id = int(data.get("id", -1)) - text = data.get("text", "") - clear_first = data.get("clear_first", True) - - with _display_lock: - user_disp = _user_displays.get(sess["user_id"], {}) - elements = user_disp.get("last_elements", []) - cached_url = user_disp.get("last_elements_url", "") - - if element_id < 0 or not elements: - await send({ - "type": "error", - "msg": "Call browser_get_elements first, then use the returned element id", - }) - return - - if await _elements_are_stale(display, cached_url): - await send({ - "type": "error", - "msg": "Page changed since last browser_get_elements — element coordinates are stale. Call browser_get_elements again before typing.", - "stale": True, - }) - return - - type_info = await cdp_type_into_element(display, element_id, elements, text, clear_first) - if "error" in type_info: - await send({"type": "error", "msg": type_info["error"]}) - return - - x, y = type_info["x"], type_info["y"] - await xdo(["mousemove", "--sync", str(x), str(y)], display) - await asyncio.sleep(0.07) - await xdo(["click", "1"], display) - await asyncio.sleep(0.05) - if clear_first: - await xdo(["key", "--clearmodifiers", "ctrl+a"], display) - await xdo(["key", "--clearmodifiers", "Delete"], display) - await type_smart(text, display) - - await send({ - "type": "ack", - "action": "browser_type_into_element", - "element_id": element_id, - "x": x, "y": y, - "text_len": len(text), - }) - asyncio.create_task(shot_bg("after typing into element", delay=0.4)) - - # ── browser_find_element ───────────────────────────────── - # بحث نصي عن عناصر بدل تخمين id يدوياً — يفترض browser_get_elements - # نُفّذ من قبل ضمن نفس الجلسة (يستخدم last_elements المخزّنة) - elif action == "browser_find_element": - query = data.get("query", "") - with _display_lock: - elements = _user_displays.get(sess["user_id"], {}).get("last_elements", []) - if not elements: - await send({ - "type": "browser_find_result", - "matches": [], - "error": "No cached elements — call browser_get_elements first", - }) - return - matches = cdp_find_elements_by_text(elements, query, limit=int(data.get("limit", 8))) - await send({ - "type": "browser_find_result", - "query": query, - "count": len(matches), - "matches": matches, - "summary": _format_elements_for_llm(matches) if matches else f"No element matching '{query}' found.", - }) - - # ── browser_get_page_text ──────────────────────────────── - # نص الصفحة الكامل — مثل Manus shell_view لكن للمتصفح - elif action == "browser_get_page_text": - result = await cdp_get_page_text(display) - await send({ - "type": "browser_page_text", - "text": result.get("text", ""), - "title": result.get("title", ""), - "url": result.get("url", ""), - "char_count": result.get("char_count", 0), - "error": result.get("error"), - }) - - # ── browser_execute_js ─────────────────────────────────── - # تنفيذ JavaScript — مثل Manus browser_console_exec - elif action == "browser_execute_js": - js_code = data.get("code", "") - if not js_code: - await send({"type": "error", "msg": "No JS code provided"}) - return - result = await cdp_execute_js(display, js_code) - await send({ - "type": "browser_js_result", - "result": result.get("result"), - "type_name": result.get("type", "undefined"), - "error": result.get("error"), - }) - - # ── browser_navigate ───────────────────────────────────── - # انتقال سريع لرابط عبر CDP بدون إعادة فتح المتصفح - elif action == "browser_navigate": - url = _safe_search(data.get("url", "")) - if not url: - await send({"type": "error", "msg": "No URL provided"}) - return - - if BROWSER_SUPPORTS_CDP: - result = await cdp_navigate(display, url) - if "ok" in result: - await send({"type": "ack", "action": "browser_navigate", "url": url, "cdp": True}) - # انتظر تحميل الصفحة ثم التقط - asyncio.create_task(shot_bg("after cdp navigate", delay=1.5, extra_shot_delay=2.5)) - else: - # fallback لـ xdotool إذا فشل CDP - await xdo(["key", "--clearmodifiers", "ctrl+l"], display) - await asyncio.sleep(0.2) - await type_smart(url, display) - await xdo(["key", "--clearmodifiers", "Return"], display) - await send({"type": "ack", "action": "browser_navigate", "url": url, "cdp": False}) - else: - await xdo(["key", "--clearmodifiers", "ctrl+l"], display) - await asyncio.sleep(0.2) - await type_smart(url, display) - await xdo(["key", "--clearmodifiers", "Return"], display) - await send({"type": "ack", "action": "browser_navigate", "url": url, "cdp": False}) - # ── تعطيل الالتقاط التلقائي — shot_bg أعلاه يتكفل بها ── - - # ── browser_scroll_to_element ──────────────────────────── - # تمرير الصفحة لعنصر معين بـ index - elif action == "browser_scroll_to_element": - element_id = int(data.get("id", -1)) - with _display_lock: - elements = _user_displays.get(sess["user_id"], {}).get("last_elements", []) - el = next((e for e in elements if e.get("id") == element_id), None) - if not el: - await send({"type": "error", "msg": f"Element {element_id} not found"}) - return - # استخدم CDP لتمرير الصفحة للعنصر - await cdp_execute_js(display, - f"document.elementFromPoint({el['x']}, {el['y']})?.scrollIntoView({{block:'center',behavior:'smooth'}})" - ) - await asyncio.sleep(0.5) - await send({"type": "ack", "action": "browser_scroll_to_element", "element_id": element_id}) - asyncio.create_task(shot_bg("after scroll to element", delay=0.8)) - - # ── browser_get_state ───────────────────────────────────── - # نداء واحد يجمع العناصر + نص الصفحة + لقطة شاشة — بدل 3 نداءات - # منفصلة. هذا أقرب لفلسفة Manus في الـ agent loop: "دورة واحدة، - # ملاحظة واحدة شاملة عن الحالة الراهنة، ثم قرار" بدل تعدد نداءات - # الملاحظة قبل كل قرار فعلي. - elif action == "browser_get_state": - elements_result = await cdp_get_elements(display) if BROWSER_SUPPORTS_CDP else {"elements": [], "error": "CDP not available"} - elements = elements_result.get("elements", []) - - with _display_lock: - if sess["user_id"] in _user_displays: - _user_displays[sess["user_id"]]["last_elements"] = elements - _user_displays[sess["user_id"]]["last_elements_url"] = elements_result.get("page_url", "") - - text_result = await cdp_get_page_text(display) if BROWSER_SUPPORTS_CDP else {"text": ""} - shot = await _priority_capture(0.6, 70) - - await send({ - "type": "browser_state", - "elements": elements, - "count": len(elements), - "page_title": elements_result.get("page_title", text_result.get("title", "")), - "page_url": elements_result.get("page_url", text_result.get("url", "")), - "viewport_w": elements_result.get("viewport_w", 1280), - "viewport_h": elements_result.get("viewport_h", 800), - "page_text_preview": (text_result.get("text", "") or "")[:2000], - "summary": _format_elements_for_llm(elements), - "screenshot": shot.get("data", "") if shot else "", - "screen_width": shot.get("width") if shot else None, - "screen_height": shot.get("height") if shot else None, - "error": elements_result.get("error"), - }) - - # ── unknown ────────────────────────────────────────────── - else: - await send({"type": "error", "msg": f"Unknown action: '{action}'"}) - - -# ════════════════════════════════════════════════════════════════ -# ── WebSocket Endpoint ──────────────────────────────────────── -# ════════════════════════════════════════════════════════════════ - -@app.websocket("/ws") -async def websocket_endpoint( - ws: WebSocket, - user_id: str = Query(default="anonymous") -): - await ws.accept() - - # الحصول على session المستخدم أو إنشاء واحدة جديدة - sess = await get_or_create_user_session(user_id, ws) - - # تأكد أن Xvfb يعمل لهذا المستخدم - await ensure_xvfb_for_session(sess) - - display = sess["display"] - - async def _heartbeat(): - while True: - await asyncio.sleep(20) - try: - await ws.send_text(json.dumps({"type": "ping", "ts": int(time.time()*1000)})) - except Exception: - break - - hb_task = asyncio.create_task(_heartbeat()) - - try: - w, h = await asyncio.to_thread(_get_screen_size, display) - await ws.send_text(json.dumps({ - "type": "connected", - "screen_width": w, "screen_height": h, - "browser": BROWSER, - "display": display, - "user_id": user_id, - "session_id": id(ws), - "msg": f"Z Computer Mode v11 | User: {user_id} | Display: {display} | Browser: {BROWSER} | Screen: {w}x{h}", - }, ensure_ascii=False)) - - # لقطة شاشة أولية - result = await asyncio.to_thread(capture_with_grid, display, 0.65, 72) - if result["data"]: - await ws.send_text(json.dumps({ - "type": "screenshot", - "data": result["data"], - "data_grid": result.get("data_grid", ""), - "data_grid2": result.get("data_grid2", ""), - "data_grid3": result.get("data_grid3", ""), - "ts": int(time.time() * 1000), - "label": "Initial screen", - "screen_width": result["width"], - "screen_height": result["height"], - "mouse_x": result["mouse_x"], - "mouse_y": result["mouse_y"], - "has_grid": True, - }, ensure_ascii=False)) - sess["last_bg_hash"] = _frame_hash(result["data"]) - - except Exception as e: - print(f"[ws:{user_id}] init error: {e}") - - try: - while True: - raw = await ws.receive_text() - try: - msg = json.loads(raw) - if msg.get("type") == "pong": - continue - await handle_action(ws, msg, sess) - except json.JSONDecodeError: - pass - except WebSocketDisconnect: - pass - except Exception as e: - print(f"[ws:{user_id}] error: {e}") - finally: - hb_task.cancel() - await destroy_user_ws(user_id, ws) - - -# ════════════════════════════════════════════════════════════════ -# ── REST Endpoints ──────────────────────────────────────────── -# ════════════════════════════════════════════════════════════════ - -@app.get("/screenshot") -async def rest_screenshot(user_id: str = "anonymous"): - with _display_lock: - sess = _user_displays.get(user_id) - display = sess["display"] if sess else f":{DISPLAY_BASE}" - result = await asyncio.to_thread(capture_with_grid, display, 0.7, 75) - return JSONResponse({ - "image": result["data"], - "image_grid": result.get("data_grid", ""), - "image_grid2": result.get("data_grid2", ""), - "image_grid3": result.get("data_grid3", ""), - "ts": int(time.time() * 1000), - "screen_width": result["width"], - "screen_height": result["height"], - "mouse_x": result["mouse_x"], - "mouse_y": result["mouse_y"], - "has_grid": True, - "display": display, - "user_id": user_id, - }) - -@app.post("/terminal") -async def rest_terminal(body: dict): - user_id = body.get("user_id", "anonymous") - with _display_lock: - sess = _user_displays.get(user_id) - display = sess["display"] if sess else f":{DISPLAY_BASE}" - return JSONResponse(await run_cmd_smart(body.get("cmd", ""), display, body.get("timeout", 60))) - -@app.get("/health") -async def health(): - with _display_lock: - n = len(_user_displays) - users = [ - { - "user_id": uid, - "display": s["display"], - "cdp_port": s.get("cdp_port"), - "connected": bool(s.get("active_ws")), - "elements_cached": len(s.get("last_elements", [])), - } - for uid, s in _user_displays.items() - ] - with _cdp_sessions_lock: - cdp_active = sum(1 for s in _cdp_sessions.values() if s._ws) - return { - "status": "ok", - "version": "v12-cdp-dom-elements", - "browser": BROWSER, - "cdp_supported": BROWSER_SUPPORTS_CDP, - "cdp_library": _CDP_AVAILABLE, - "active_users": n, - "cdp_sessions_active": cdp_active, - "users": users, - } - - -# ════════════════════════════════════════════════════════════════ -# ── Background Tasks ────────────────────────────────────────── -# ════════════════════════════════════════════════════════════════ - -async def _cleanup_tmp(): - """ينظّف /tmp كل 5 دقائق. ملاحظة v13: محرك الالتقاط - الجديد (import-pipe/xlib-direct/xwd) لا يكتب أي ملفات zss_* على الإطلاق (كلشيء - في الذاكرة عبر pipe)، فهذه الدالة أصبحت غير ضرورية فعلياً للقطات - الشاشة، لكنها تُرك كشبكة أمان إضافية (مثلاً لو أضيف كود مستقبلاً - يكتب ملفات مؤقتة بنفس البادئة).""" - while True: - await asyncio.sleep(300) - try: - subprocess.run( - ["find", "/tmp", "-name", "zss_*", "-mmin", "+10", "-delete"], - capture_output=True, timeout=10 - ) - except Exception as e: - print(f"[cleanup] {e}") - - -async def _cleanup_idle_sessions(): - """ - يُزيل sessions المستخدمين غير النشطين (لا اتصال منذ أكثر من ساعة) - لتحرير الـ displays والذاكرة. - """ - while True: - await asyncio.sleep(1800) # كل 30 دقيقة - now = time.time() - to_remove = [] - with _display_lock: - for uid, sess in list(_user_displays.items()): - if sess.get("active_ws"): - continue # لا تحذف sessions النشطة - if now - sess.get("created", now) > 3600: # أكثر من ساعة - to_remove.append((uid, sess)) - - for uid, sess in to_remove: - print(f"[cleanup] 🗑️ Removing idle session for '{uid}' on {sess['display']}") - _kill_proc(sess.get("browser_proc")) - _kill_proc(sess.get("xvfb_proc")) - _kill_display_processes(sess["display"]) - with _display_lock: - _user_displays.pop(uid, None) - try: - _display_numbers.discard(int(sess["display"].lstrip(":"))) - except Exception: - pass - - -@app.on_event("startup") -async def startup(): - asyncio.create_task(_cleanup_tmp()) - asyncio.create_task(_cleanup_idle_sessions()) - print("✅ Z Computer Mode v12 ready — CDP DOM Elements + Per-User Isolated Displays") - print(f" Display range: :{DISPLAY_BASE} → :{DISPLAY_MAX}") - print(f" Browser: {BROWSER}") - print(f" CDP (Manus-style DOM): {'✅ ENABLED' if BROWSER_SUPPORTS_CDP else '⚠️ DISABLED (Firefox)'}") - print(f" websockets lib: {'✅' if _CDP_AVAILABLE else '❌ (pip install websockets)'}") - print(f" CDP port range: {CDP_PORT_BASE} → {CDP_PORT_BASE + (DISPLAY_MAX - DISPLAY_BASE)}") - print(f" Connect: wss://your-space.hf.space/ws?user_id=EMAIL") - print() - print(" 🆕 New v12 actions:") - print(" browser_get_elements → DOM elements list with exact coords (like Manus)") - print(" browser_click_element → click by index id (no coord guessing)") - print(" browser_type_into_element → type into an input by index (click+type combined)") - print(" browser_find_element → find elements by text/description instead of scanning ids") - print(" browser_get_page_text → full page text content") - print(" browser_get_state → elements + page text + screenshot in one round trip") - print(" browser_execute_js → run JS in browser (like Manus console_exec)") - print(" browser_navigate → fast CDP navigation") - print(" browser_scroll_to_element → scroll to element by id") - - -if __name__ == "__main__": - port = int(os.environ.get("PORT", 7860)) - uvicorn.run("app:app", host="0.0.0.0", port=port, log_level="info")