Spaces:
Running
Running
Download Browser_Agent/job_queue/task_queue.py from Luca448/APP-Backend: direct link, hf CLI and curl.
- Browser
- Download file 39.8 kB
-
https://huggingface.co/spaces/Luca448/APP-Backend/resolve/main/Browser_Agent/job_queue/task_queue.py
- Command line
-
hf download hf://spaces/Luca448/APP-Backend/Browser_Agent/job_queue/task_queue.py
-
curl -L -o task_queue.py https://huggingface.co/spaces/Luca448/APP-Backend/resolve/main/Browser_Agent/job_queue/task_queue.py
39.8 kB
| # queue/task_queue.py | |
| import asyncio | |
| import logging | |
| import uuid | |
| import hashlib | |
| import json | |
| from dataclasses import dataclass, field | |
| from datetime import datetime, timedelta, timezone | |
| from enum import Enum | |
| from config.settings import settings | |
| from security.crypto import encrypt_state, decrypt_state | |
| from job_queue.task_store import TaskStore, task_store | |
| logger = logging.getLogger(__name__) | |
| class TaskStatus(str, Enum): | |
| PENDING = "pending" | |
| PLANNING = "planning" | |
| RUNNING = "running" | |
| RETRYING = "retrying" | |
| WAITING_FOR_CONFIRMATION = "waiting_for_confirmation" | |
| WAITING_FOR_HANDOFF = "waiting_for_handoff" | |
| DONE = "done" | |
| FAILED = "failed" | |
| CANCELLED = "cancelled" | |
| TIMED_OUT = "timed_out" | |
| class Task: | |
| id: str | |
| task_type: str | |
| params: dict | |
| status: TaskStatus = TaskStatus.PENDING | |
| created_at: datetime = field(default_factory=datetime.utcnow) | |
| started_at: datetime | None = None | |
| finished_at: datetime | None = None | |
| result: str | None = None | |
| error: str | None = None | |
| chat_id: int | None = None | |
| connection_mode: str = "normal" | |
| profile_state: str | None = None | |
| save_profile: bool = False | |
| steps: list[dict] = field(default_factory=list) | |
| events: list[dict] = field(default_factory=list) | |
| revision: int = 0 | |
| idempotency_key: str | None = None | |
| request_digest: str | None = None | |
| deadline_at: datetime | None = None | |
| heartbeat_at: datetime | None = None | |
| handoff: dict | None = None | |
| confirmation: dict | None = None | |
| error_details: dict | None = None | |
| attempt: int = 0 | |
| cancel_requested: bool = False | |
| def add_event( | |
| self, | |
| event_type: str, | |
| message: str, | |
| *, | |
| status: str | None = None, | |
| details: dict | None = None, | |
| ) -> dict: | |
| self.revision += 1 | |
| event = { | |
| "sequence": len(self.events) + 1, | |
| "timestamp": datetime.now(timezone.utc).isoformat(), | |
| "type": event_type, | |
| "status": status or self.status.value, | |
| "message": message, | |
| "details": details or {}, | |
| } | |
| self.events.append(event) | |
| return event | |
| def is_terminal(self) -> bool: | |
| return self.status in { | |
| TaskStatus.DONE, | |
| TaskStatus.FAILED, | |
| TaskStatus.CANCELLED, | |
| TaskStatus.TIMED_OUT, | |
| } | |
| class TaskQueue: | |
| """ | |
| Asynchrone In-Process-Task-Queue. | |
| """ | |
| def __init__(self, max_size: int = 20, store: TaskStore | None = None): | |
| self._queue: asyncio.Queue[Task] = asyncio.Queue(maxsize=max_size) | |
| self._tasks: dict[str, Task] = {} | |
| self._running_tasks: dict[str, asyncio.Task] = {} | |
| self._max_size = max_size | |
| self._store = store or task_store | |
| self._idempotency: dict[str, tuple[str, str]] = {} | |
| self.worker_heartbeat: datetime | None = None | |
| def _request_digest(task_type: str, params: dict, chat_id: int | None) -> str: | |
| raw = json.dumps( | |
| {"task_type": task_type, "params": params, "chat_id": chat_id}, | |
| sort_keys=True, | |
| ensure_ascii=False, | |
| default=str, | |
| ) | |
| return hashlib.sha256(raw.encode("utf-8")).hexdigest() | |
| def _store_payload(self, task: Task) -> dict: | |
| sensitive_keys = { | |
| "password", | |
| "cookies", | |
| "cookie", | |
| "token", | |
| "api_key", | |
| "encrypted_state", | |
| "form_data", | |
| } | |
| def redact(value): | |
| if isinstance(value, dict): | |
| return { | |
| str(key): "[REDACTED]" if str(key).lower() in sensitive_keys else redact(item) | |
| for key, item in value.items() | |
| } | |
| if isinstance(value, list): | |
| return [redact(item) for item in value] | |
| return value | |
| public_handoff = None | |
| if task.handoff: | |
| public_handoff = { | |
| key: value | |
| for key, value in task.handoff.items() | |
| if key not in {"cookies", "form_data"} | |
| } | |
| public_payload = { | |
| "result": None, | |
| "error": None, | |
| "error_details": task.error_details, | |
| "steps": redact(task.steps), | |
| "events": redact(task.events), | |
| "handoff": public_handoff, | |
| "confirmation": task.confirmation, | |
| "connection_mode": task.connection_mode, | |
| "save_profile": task.save_profile, | |
| } | |
| private_payload = { | |
| "params": task.params, | |
| "profile_state": task.profile_state, | |
| "chat_id": task.chat_id, | |
| "result": task.result, | |
| "error": task.error, | |
| "handoff": task.handoff, | |
| "steps": task.steps, | |
| "events": task.events, | |
| } | |
| return { | |
| "id": task.id, | |
| "task_type": task.task_type, | |
| "status": task.status.value, | |
| "revision": task.revision, | |
| "idempotency_key": task.idempotency_key, | |
| "params_encrypted": encrypt_state(private_payload), | |
| "public_payload": public_payload, | |
| "created_at": task.created_at.replace(tzinfo=timezone.utc).isoformat(), | |
| "started_at": task.started_at.replace(tzinfo=timezone.utc).isoformat() if task.started_at else None, | |
| "finished_at": task.finished_at.replace(tzinfo=timezone.utc).isoformat() if task.finished_at else None, | |
| "deadline_at": task.deadline_at.isoformat() if task.deadline_at else None, | |
| "heartbeat_at": task.heartbeat_at.isoformat() if task.heartbeat_at else None, | |
| "updated_at": datetime.now(timezone.utc).isoformat(), | |
| } | |
| async def persist(self, task: Task, event: dict | None = None) -> None: | |
| await self._store.save(self._store_payload(task)) | |
| if event: | |
| await self._store.append_event({ | |
| "task_id": task.id, | |
| "sequence": event["sequence"], | |
| "event_type": event["type"], | |
| "payload": event, | |
| }) | |
| async def initialize(self) -> None: | |
| await self._store.initialize() | |
| for row in await self._store.load_active(): | |
| task_id = str(row.get("id", "")) | |
| if not task_id or task_id in self._tasks: | |
| continue | |
| private = decrypt_state(str(row.get("params_encrypted") or "")) or {} | |
| public = row.get("public_payload") or {} | |
| try: | |
| status = TaskStatus(str(row.get("status", TaskStatus.PENDING.value))) | |
| except ValueError: | |
| status = TaskStatus.FAILED | |
| task = Task( | |
| id=task_id, | |
| task_type=str(row.get("task_type", "agent_task")), | |
| params=dict(private.get("params") or {}), | |
| status=status, | |
| created_at=datetime.fromisoformat(str(row["created_at"])).replace(tzinfo=None), | |
| result=private.get("result"), | |
| error=private.get("error"), | |
| chat_id=private.get("chat_id"), | |
| connection_mode=str(public.get("connection_mode", "normal")), | |
| profile_state=private.get("profile_state"), | |
| save_profile=bool(public.get("save_profile", False)), | |
| steps=list(private.get("steps") or public.get("steps") or []), | |
| events=list(private.get("events") or public.get("events") or []), | |
| revision=int(row.get("revision") or 0), | |
| idempotency_key=row.get("idempotency_key"), | |
| handoff=private.get("handoff") or public.get("handoff"), | |
| confirmation=public.get("confirmation"), | |
| error_details=public.get("error_details"), | |
| ) | |
| self._tasks[task.id] = task | |
| if task.idempotency_key: | |
| digest = self._request_digest(task.task_type, task.params, task.chat_id) | |
| self._idempotency[task.idempotency_key] = (task.id, digest) | |
| if status in {TaskStatus.PENDING, TaskStatus.PLANNING, TaskStatus.RETRYING}: | |
| task.status = TaskStatus.PENDING | |
| await self._queue.put(task) | |
| elif status == TaskStatus.RUNNING: | |
| if task.params.get("safe_to_retry") is True: | |
| task.status = TaskStatus.PENDING | |
| await self._queue.put(task) | |
| else: | |
| task.status = TaskStatus.FAILED | |
| task.error = "Der Server wurde während einer nicht sicher wiederholbaren Aktion neu gestartet." | |
| task.error_details = { | |
| "code": "restart_recovery_required", | |
| "retryable": True, | |
| "user_action": "Aufgabe erneut starten", | |
| } | |
| task.finished_at = datetime.utcnow() | |
| await self.persist(task) | |
| async def enqueue( | |
| self, | |
| task_type: str, | |
| params: dict, | |
| chat_id: int | None = None, | |
| timeout: float = 5.0, | |
| idempotency_key: str | None = None, | |
| ) -> str: | |
| digest = self._request_digest(task_type, params, chat_id) | |
| if idempotency_key and idempotency_key in self._idempotency: | |
| existing_id, existing_digest = self._idempotency[idempotency_key] | |
| if existing_digest != digest: | |
| raise ValueError("idempotency_conflict") | |
| return existing_id | |
| task_id = f"task_{uuid.uuid4().hex}" | |
| task = Task( | |
| id=task_id, | |
| task_type=task_type, | |
| params=params, | |
| chat_id=chat_id, | |
| connection_mode=params.get("connection_mode", "normal"), | |
| profile_state=params.get("profile_state"), | |
| save_profile=params.get("save_profile", False), | |
| idempotency_key=idempotency_key, | |
| request_digest=digest, | |
| deadline_at=datetime.now(timezone.utc) + timedelta( | |
| seconds=settings.task_active_timeout_seconds | |
| ), | |
| ) | |
| event = task.add_event("task_queued", "Alles klar, ich lege los.", status="pending") | |
| # Register and persist before the task becomes visible to a worker. | |
| self._tasks[task_id] = task | |
| if idempotency_key: | |
| self._idempotency[idempotency_key] = (task_id, digest) | |
| await self.persist(task, event) | |
| try: | |
| await asyncio.wait_for(self._queue.put(task), timeout=timeout) | |
| except (asyncio.TimeoutError, asyncio.QueueFull): | |
| self._tasks.pop(task_id, None) | |
| if idempotency_key: | |
| self._idempotency.pop(idempotency_key, None) | |
| raise asyncio.QueueFull(f"Queue voll (max {self._max_size} Tasks). Bitte warten.") | |
| logger.info(f"Task eingereiht: {task_id} (Typ: {task_type})") | |
| return task_id | |
| async def get_status(self, task_id: str, *, typed: bool = False) -> dict: | |
| task = self._tasks.get(task_id) | |
| if not task: | |
| return {"error": f"Task '{task_id}' nicht gefunden"} | |
| await self.persist(task) | |
| public_handoff = None | |
| if task.handoff: | |
| public_handoff = { | |
| key: value for key, value in task.handoff.items() | |
| if key not in {"cookies", "form_data"} | |
| } | |
| public_handoff["payload_available"] = True | |
| status = task.status.value | |
| if typed: | |
| status = { | |
| "pending": "queued", | |
| "done": "succeeded", | |
| }.get(status, status) | |
| return { | |
| "id": task.id, | |
| "status": status, | |
| "task_type": task.task_type, | |
| "created_at": task.created_at.isoformat(), | |
| "started_at": task.started_at.isoformat() if task.started_at else None, | |
| "finished_at": task.finished_at.isoformat() if task.finished_at else None, | |
| "result": task.result, | |
| "error": task.error, | |
| "error_details": task.error_details, | |
| "steps": task.steps, | |
| "events": task.events, | |
| "revision": task.revision, | |
| "handoff": public_handoff, | |
| "confirmation": task.confirmation, | |
| } | |
| async def get_handoff_payload(self, task_id: str, handoff_id: str) -> dict | None: | |
| task = self._tasks.get(task_id) | |
| if not task or not task.handoff or task.handoff.get("id") != handoff_id: | |
| return None | |
| expires_at = datetime.fromisoformat(str(task.handoff["expires_at"])) | |
| if expires_at <= datetime.now(timezone.utc): | |
| return None | |
| return { | |
| "id": handoff_id, | |
| "url": task.handoff.get("url"), | |
| "cookies": task.handoff.get("cookies") or [], | |
| "form_data": task.handoff.get("form_data") or {}, | |
| "completion_policy": task.handoff.get("completion_policy"), | |
| "expires_at": task.handoff.get("expires_at"), | |
| "revision": task.revision, | |
| } | |
| async def resume_handoff( | |
| self, | |
| task_id: str, | |
| handoff_id: str, | |
| *, | |
| revision: int, | |
| cookies: list, | |
| resume_url: str, | |
| submission_confirmed: bool, | |
| ) -> Task: | |
| task = self._tasks.get(task_id) | |
| if task and task.params.get("_completed_handoff_id") == handoff_id: | |
| return task | |
| if not task or not task.handoff or task.handoff.get("id") != handoff_id: | |
| raise KeyError("handoff_not_found") | |
| if task.status != TaskStatus.WAITING_FOR_HANDOFF: | |
| raise ValueError("handoff_not_waiting") | |
| expires_at = datetime.fromisoformat(str(task.handoff["expires_at"])) | |
| if expires_at <= datetime.now(timezone.utc): | |
| raise ValueError("handoff_expired") | |
| if revision != task.revision: | |
| raise ValueError("revision_conflict") | |
| if task.handoff.get("completion_policy") == "registration_submitted" and not submission_confirmed: | |
| raise ValueError("submission_not_confirmed") | |
| task.params["exported_cookies"] = cookies | |
| task.params["_resume_handoff"] = { | |
| "id": handoff_id, | |
| "resume_url": resume_url, | |
| "submission_confirmed": submission_confirmed, | |
| } | |
| task.params["_completed_handoff_id"] = handoff_id | |
| task.handoff = None | |
| task.status = TaskStatus.PENDING | |
| task.finished_at = None | |
| task.deadline_at = datetime.now(timezone.utc) + timedelta( | |
| seconds=settings.task_active_timeout_seconds | |
| ) | |
| event = task.add_event("handoff_completed", "Danke, ich mache direkt weiter.", status="pending") | |
| await self.persist(task, event) | |
| await self._queue.put(task) | |
| return task | |
| async def resolve_confirmation( | |
| self, | |
| task_id: str, | |
| *, | |
| confirmation_id: str, | |
| revision: int, | |
| approved: bool, | |
| ) -> Task: | |
| task = self._tasks.get(task_id) | |
| if not task or not task.confirmation or task.confirmation.get("id") != confirmation_id: | |
| raise KeyError("confirmation_not_found") | |
| if task.status != TaskStatus.WAITING_FOR_CONFIRMATION: | |
| raise ValueError("confirmation_not_waiting") | |
| if revision != task.revision: | |
| raise ValueError("revision_conflict") | |
| expires_at = datetime.fromisoformat(str(task.confirmation["expires_at"])) | |
| if expires_at <= datetime.now(timezone.utc): | |
| raise ValueError("confirmation_expired") | |
| if not approved: | |
| task.status = TaskStatus.CANCELLED | |
| task.result = "Aktion nicht ausgeführt – Bestätigung wurde abgelehnt." | |
| task.finished_at = datetime.utcnow() | |
| task.confirmation = None | |
| event = task.add_event( | |
| "confirmation_declined", | |
| "Verstanden, ich führe diese Aktion nicht aus.", | |
| status="cancelled", | |
| ) | |
| await self.persist(task, event) | |
| return task | |
| checkpoint = task.params.pop("_confirmation_checkpoint", None) | |
| if not isinstance(checkpoint, dict) or checkpoint.get("id") != confirmation_id: | |
| raise ValueError("confirmation_checkpoint_missing") | |
| task.params["exported_cookies"] = checkpoint.get("cookies") or [] | |
| task.params["_approved_confirmation"] = { | |
| **task.confirmation, | |
| "consumed": False, | |
| } | |
| task.params["_resume_confirmation"] = { | |
| "id": confirmation_id, | |
| "url": checkpoint.get("url"), | |
| } | |
| task.confirmation = None | |
| task.status = TaskStatus.PENDING | |
| task.finished_at = None | |
| task.deadline_at = datetime.now(timezone.utc) + timedelta( | |
| seconds=settings.task_active_timeout_seconds | |
| ) | |
| event = task.add_event( | |
| "confirmation_approved", | |
| "Bestätigt – ich führe genau diese Aktion jetzt aus.", | |
| status="pending", | |
| ) | |
| await self.persist(task, event) | |
| await self._queue.put(task) | |
| return task | |
| async def get_next(self) -> Task: | |
| return await self._queue.get() | |
| async def cancel(self, task_id: str) -> bool: | |
| task = self._tasks.get(task_id) | |
| if not task: | |
| return False | |
| if task.is_terminal: | |
| return True | |
| task.cancel_requested = True | |
| task.status = TaskStatus.CANCELLED | |
| task.finished_at = datetime.utcnow() | |
| task.error = None | |
| task.result = "Aufgabe vom Nutzer gestoppt." | |
| event = task.add_event("task_cancelled", "Okay, ich habe die Aufgabe gestoppt.", status="cancelled") | |
| running = self._running_tasks.get(task_id) | |
| if running and not running.done(): | |
| running.cancel() | |
| logger.info("Task %s wurde vom Nutzer gestoppt", task_id) | |
| await self.persist(task, event) | |
| return True | |
| def task_done(self): | |
| self._queue.task_done() | |
| def size(self) -> int: | |
| return self._queue.qsize() | |
| def list_tasks(self) -> list[dict]: | |
| return [ | |
| { | |
| "id": t.id, | |
| "type": t.task_type, | |
| "status": t.status, | |
| "created": t.created_at.isoformat(), | |
| } | |
| for t in self._tasks.values() | |
| ] | |
| task_queue = TaskQueue() | |
| async def worker_loop(queue: TaskQueue, brain, browser_mgr, telegram_bot=None): | |
| """ | |
| Der Worker-Loop: Verarbeitet Aufgaben aus der Queue sequenziell. | |
| """ | |
| logger.info("Worker-Loop gestartet — warte auf Aufgaben...") | |
| while True: | |
| queue.worker_heartbeat = datetime.now(timezone.utc) | |
| task = await queue.get_next() | |
| if task.cancel_requested or task.status == TaskStatus.CANCELLED: | |
| queue.task_done() | |
| continue | |
| task.status = TaskStatus.RUNNING | |
| task.started_at = task.started_at or datetime.utcnow() | |
| task.heartbeat_at = datetime.now(timezone.utc) | |
| task.attempt += 1 | |
| started_event = task.add_event( | |
| "task_started", | |
| "Ich schaue mir die Aufgabe jetzt genauer an.", | |
| status="running", | |
| ) | |
| await queue.persist(task, started_event) | |
| logger.info(f"Verarbeite Task {task.id}: {task.task_type}") | |
| try: | |
| execution = asyncio.create_task(_dispatch_task(task, brain, browser_mgr)) | |
| queue._running_tasks[task.id] = execution | |
| remaining = settings.task_active_timeout_seconds | |
| if task.deadline_at: | |
| remaining = max( | |
| 0.1, | |
| (task.deadline_at - datetime.now(timezone.utc)).total_seconds(), | |
| ) | |
| result = await asyncio.wait_for(execution, timeout=remaining) | |
| if not task.cancel_requested: | |
| if isinstance(result, dict) and result.get("status") == "handoff_required": | |
| task.status = TaskStatus.WAITING_FOR_HANDOFF | |
| task.handoff = dict(result.get("handoff") or {}) | |
| task.result = result.get("handoff_signal") | |
| task.finished_at = None | |
| task.deadline_at = None | |
| event = task.add_event( | |
| "handoff_required", | |
| str(result.get("message") or "Ich brauche kurz deine Hilfe im Browser."), | |
| status="waiting_for_handoff", | |
| details={ | |
| "handoff_id": task.handoff.get("id"), | |
| "type": task.handoff.get("type"), | |
| }, | |
| ) | |
| await queue.persist(task, event) | |
| elif isinstance(result, dict) and result.get("status") == "confirmation_required": | |
| task.status = TaskStatus.WAITING_FOR_CONFIRMATION | |
| task.confirmation = dict(result.get("confirmation") or {}) | |
| task.result = None | |
| task.finished_at = None | |
| task.deadline_at = None | |
| event = task.add_event( | |
| "confirmation_required", | |
| str(result.get("message") or "Bevor ich weitermache, brauche ich deine Bestätigung."), | |
| status="waiting_for_confirmation", | |
| ) | |
| await queue.persist(task, event) | |
| elif isinstance(result, dict) and ( | |
| result.get("status") == "error" or result.get("success") is False | |
| ): | |
| task.status = TaskStatus.FAILED | |
| task.error = str(result.get("error") or "Die Aufgabe ist fehlgeschlagen.") | |
| task.error_details = { | |
| "code": str(result.get("error_type") or "task_failed"), | |
| "retryable": bool(result.get("retryable", False)), | |
| "user_action": result.get("suggestion"), | |
| "diagnostic_id": uuid.uuid4().hex[:12], | |
| } | |
| else: | |
| task.status = TaskStatus.DONE | |
| task.result = str(result) | |
| event = task.add_event( | |
| "task_succeeded", | |
| "Fertig – ich habe die Aufgabe abgeschlossen.", | |
| status="done", | |
| ) | |
| await queue.persist(task, event) | |
| logger.info("Task %s erfolgreich", task.id) | |
| except asyncio.CancelledError: | |
| if task.status != TaskStatus.CANCELLED: | |
| task.status = TaskStatus.CANCELLED | |
| task.result = "Aufgabe vom Nutzer gestoppt." | |
| task.error = None | |
| logger.info("Task %s sauber abgebrochen", task.id) | |
| except asyncio.TimeoutError: | |
| task.status = TaskStatus.TIMED_OUT | |
| task.error = "Die Aufgabe hat das aktive Zeitlimit überschritten." | |
| task.error_details = { | |
| "code": "task_timed_out", | |
| "retryable": True, | |
| "user_action": "Aufgabe erneut versuchen", | |
| "diagnostic_id": uuid.uuid4().hex[:12], | |
| } | |
| task.add_event( | |
| "task_timed_out", | |
| "Das dauert ungewöhnlich lange. Ich habe sicher abgebrochen.", | |
| status="timed_out", | |
| ) | |
| except Exception as e: | |
| task.status = TaskStatus.FAILED | |
| task.error = str(e) | |
| task.error_details = { | |
| "code": "unhandled_task_error", | |
| "retryable": False, | |
| "user_action": None, | |
| "diagnostic_id": uuid.uuid4().hex[:12], | |
| } | |
| task.add_event( | |
| "task_failed", | |
| "Da ist etwas schiefgelaufen. Die Diagnose ist gespeichert.", | |
| status="failed", | |
| ) | |
| logger.error(f"Task {task.id} fehlgeschlagen: {e}", exc_info=True) | |
| finally: | |
| if task.status not in { | |
| TaskStatus.WAITING_FOR_HANDOFF, | |
| TaskStatus.WAITING_FOR_CONFIRMATION, | |
| }: | |
| task.finished_at = datetime.utcnow() | |
| task.heartbeat_at = datetime.now(timezone.utc) | |
| queue._running_tasks.pop(task.id, None) | |
| queue.task_done() | |
| await queue.persist(task) | |
| if telegram_bot and task.chat_id: | |
| try: | |
| await _notify_telegram(telegram_bot, task) | |
| except Exception as e: | |
| logger.warning(f"Telegram-Benachrichtigung fehlgeschlagen: {e}") | |
| async def _dispatch_task(task: Task, brain, browser_mgr) -> str: | |
| decrypted_state = None | |
| if task.profile_state: | |
| decrypted_state = decrypt_state(task.profile_state) | |
| elif task.params.get("exported_cookies"): | |
| raw_cookies = task.params.get("exported_cookies", []) | |
| formatted_cookies = [] | |
| for c in raw_cookies: | |
| if not isinstance(c, dict) or not c.get("name"): | |
| continue | |
| cd = { | |
| "name": str(c.get("name", "")), | |
| "value": str(c.get("value", "")), | |
| "domain": str(c.get("domain", "")), | |
| "path": str(c.get("path", "/")), | |
| } | |
| if c.get("expiresDate") is not None: | |
| try: | |
| exp = float(c["expiresDate"]) | |
| cd["expires"] = exp / 1000.0 if exp > 10000000000 else exp | |
| except Exception: | |
| pass | |
| elif c.get("expires") is not None: | |
| try: | |
| cd["expires"] = float(c["expires"]) | |
| except Exception: | |
| pass | |
| if "isHttpOnly" in c: | |
| cd["httpOnly"] = bool(c["isHttpOnly"]) | |
| elif "httpOnly" in c: | |
| cd["httpOnly"] = bool(c["httpOnly"]) | |
| if "isSecure" in c: | |
| cd["secure"] = bool(c["isSecure"]) | |
| elif "secure" in c: | |
| cd["secure"] = bool(c["secure"]) | |
| formatted_cookies.append(cd) | |
| decrypted_state = {"cookies": formatted_cookies, "origins": []} | |
| logger.info(f"Loaded {len(formatted_cookies)} exported cookies into ephemeral task context.") | |
| ctx = await browser_mgr.get_ephemeral_context(connection_mode=task.connection_mode, state_data=decrypted_state) | |
| try: | |
| page = await ctx.new_page() | |
| async def handle_download(download): | |
| if task.connection_mode == "ghost": | |
| # Secure Proxy Download | |
| file_path = settings.downloads_dir / f"{uuid.uuid4().hex}_{download.suggested_filename}" | |
| await download.save_as(str(file_path)) | |
| logger.info(f"Ghost Mode: File downloaded securely to {file_path}") | |
| task.steps.append({"type": "download_ready", "path": str(file_path), "filename": download.suggested_filename}) | |
| else: | |
| # Fast Stealth / Normal Mode: Handoff to frontend | |
| url = download.url | |
| cookies = await ctx.cookies() | |
| logger.info(f"Fast Mode: Handoff download URL to frontend: {url}") | |
| await download.cancel() | |
| task.steps.append({"type": "download_handoff", "url": url, "cookies": cookies, "filename": download.suggested_filename}) | |
| page.on("download", handle_download) | |
| resume_confirmation = task.params.pop("_resume_confirmation", None) | |
| if isinstance(resume_confirmation, dict): | |
| from security.ssrf_guard import validate_url | |
| confirmation_url = validate_url(str(resume_confirmation.get("url") or "")) | |
| await page.goto( | |
| confirmation_url, | |
| wait_until="domcontentloaded", | |
| timeout=30000, | |
| ) | |
| confirmation_id = str(resume_confirmation.get("id") or "") | |
| original_instruction = str(task.params.get("task", "")) | |
| continuation_note = ( | |
| "\n\n[BESTÄTIGTE AKTION FORTSETZEN]\n" | |
| f"Der Nutzer hat genau eine externe Aktion bestätigt. Verwende beim passenden " | |
| f"Submit-Klick `confirmation_id={confirmation_id}`. Beobachte die Seite vorher " | |
| "neu und führe keine andere externe Aktion mit dieser Bestätigung aus." | |
| ) | |
| model = task.params.get("model") | |
| if model and model.startswith("openrouter-"): | |
| model = model.replace("openrouter-", "", 1) | |
| return await brain.run( | |
| original_instruction + continuation_note, | |
| page, | |
| model=model, | |
| task=task, | |
| ) | |
| resume_handoff = task.params.pop("_resume_handoff", None) | |
| if isinstance(resume_handoff, dict): | |
| resume_url = str(resume_handoff.get("resume_url") or "") | |
| submission_confirmed = resume_handoff.get("submission_confirmed") is True | |
| continuation_context = task.params.get("_account_creation_context") | |
| if isinstance(continuation_context, dict): | |
| if not submission_confirmed: | |
| return { | |
| "status": "error", | |
| "success": False, | |
| "error_type": "submission_not_confirmed", | |
| "error": "Das Registrierungsformular wurde noch nicht nachweislich abgesendet.", | |
| "retryable": True, | |
| } | |
| # Do not navigate back through the challenge. The mobile user | |
| # already submitted; continue deterministically at mailbox | |
| # verification using exactly the same account context. | |
| from workflows.account_creator import AccountCreatorWorkflow | |
| workflow = AccountCreatorWorkflow(page) | |
| workflow_result = await workflow.run( | |
| str(continuation_context.get("target_url", resume_url)), | |
| task=task, | |
| existing_context=continuation_context, | |
| resume_url=resume_url, | |
| submission_already_completed=True, | |
| ) | |
| if workflow_result.get("status") == "handoff_required": | |
| return workflow_result | |
| if workflow_result.get("success"): | |
| return ( | |
| "Registrierung erfolgreich abgeschlossen.\n\n" | |
| f"[ACCOUNT_CREATED|{continuation_context.get('target_url', resume_url)}|" | |
| f"{workflow_result.get('email', continuation_context.get('email', ''))}|" | |
| f"{workflow_result.get('password', continuation_context.get('password', ''))}]" | |
| ) | |
| return workflow_result | |
| from security.ssrf_guard import validate_url | |
| safe_resume_url = validate_url(resume_url) | |
| await page.goto(safe_resume_url, wait_until="domcontentloaded", timeout=30000) | |
| original_instruction = str(task.params.get("task", "")) | |
| continuation_note = ( | |
| "\n\n[HANDOFF FORTSETZEN]\n" | |
| "Der Nutzer hat die Sicherheitsprüfung abgeschlossen. Beobachte die aktuelle " | |
| "Seite neu, wiederhole keine bereits ausgeführte externe Aktion und setze die " | |
| "ursprüngliche Aufgabe fort." | |
| ) | |
| model = task.params.get("model") | |
| if model and model.startswith("openrouter-"): | |
| model = model.replace("openrouter-", "", 1) | |
| return await brain.run( | |
| original_instruction + continuation_note, | |
| page, | |
| model=model, | |
| task=task, | |
| ) | |
| match task.task_type: | |
| case "agent_task": | |
| model = task.params.get("model") | |
| if model: | |
| if model.startswith("openrouter-"): | |
| model = model.replace("openrouter-", "", 1) | |
| # Ensure OpenRouter provider prefixes exist | |
| if "gemma" in model and not model.startswith("google/"): | |
| model = f"google/{model}" | |
| elif "llama" in model and not model.startswith("meta-llama/"): | |
| model = f"meta-llama/{model}" | |
| # Default to free versions if not specified, since user has no credits | |
| if "gemma-4" in model and not model.endswith(":free"): | |
| model = f"{model}:free" | |
| return await brain.run(task.params["task"], page, model=model, task=task) | |
| case "account_creation": | |
| from workflows.account_creator import AccountCreatorWorkflow | |
| wf = AccountCreatorWorkflow(page) | |
| result = await wf.run(task.params["target_url"], task=task) | |
| return result | |
| case "captcha_resume": | |
| from security.ssrf_guard import validate_url | |
| resume_url = validate_url(str(task.params.get("resume_url", ""))) | |
| source_task_id = str(task.params.get("source_task_id", "")).strip() | |
| submission_confirmed = task.params.get("submission_confirmed") is True | |
| source_task = task_queue._tasks.get(source_task_id) | |
| if not source_task or source_task.task_type != "agent_task": | |
| raise ValueError("Der ursprüngliche Agent-Auftrag für den CAPTCHA-Handoff wurde nicht gefunden.") | |
| task.steps.append({ | |
| "title": "🤝 CAPTCHA-Handoff übernehmen", | |
| "subtitle": "Stelle die vom Smartphone bestätigte Browser-Sitzung wieder her...", | |
| "status": "active", | |
| "time": datetime.now().strftime("%H:%M:%S"), | |
| }) | |
| await page.goto(resume_url, wait_until="domcontentloaded", timeout=30000) | |
| task.steps[-1]["status"] = "done" | |
| task.steps[-1]["subtitle"] = "Browser-Sitzung wiederhergestellt – Auftrag wird fortgesetzt" | |
| # Restore the original task configuration (especially model, | |
| # API key and any ephemeral registration context) while using | |
| # the cookies that were just exported from the user's phone. | |
| source_params = dict(source_task.params) | |
| source_params["exported_cookies"] = task.params.get("exported_cookies", []) | |
| source_params["connection_mode"] = task.connection_mode | |
| task.params = source_params | |
| original_instruction = str(source_params.get("task", "")) | |
| if not original_instruction: | |
| raise ValueError("Der ursprüngliche Agent-Auftrag enthält keine Anweisung.") | |
| continuation_context = source_params.get("_account_creation_context") | |
| if isinstance(continuation_context, dict): | |
| from workflows.account_creator import AccountCreatorWorkflow | |
| if not submission_confirmed: | |
| raise ValueError( | |
| "Das Registrierungsformular wurde im Browser-Handoff noch nicht " | |
| "nachweislich abgesendet. Bitte Formular, Bedingungen und CAPTCHA " | |
| "vollständig abschließen und danach erneut auf Fertig tippen." | |
| ) | |
| logger.info( | |
| "CAPTCHA-Handoff setzt bestätigte Registrierung deterministisch fort; " | |
| "kein erneutes Formular und kein weiterer LLM-Agent-Loop nötig." | |
| ) | |
| workflow = AccountCreatorWorkflow(page) | |
| workflow_result = await workflow.run( | |
| str(continuation_context.get("target_url", resume_url)), | |
| task=task, | |
| existing_context=continuation_context, | |
| resume_url=resume_url, | |
| submission_already_completed=True, | |
| ) | |
| if workflow_result.get("status") == "handoff_required": | |
| return workflow_result.get( | |
| "handoff_signal", | |
| workflow_result.get("message", "Browser-Handoff erforderlich."), | |
| ) | |
| if workflow_result.get("success"): | |
| return ( | |
| "Registrierung erfolgreich abgeschlossen.\n\n" | |
| f"[ACCOUNT_CREATED|{continuation_context.get('target_url', resume_url)}|" | |
| f"{workflow_result.get('email', continuation_context.get('email', ''))}|" | |
| f"{workflow_result.get('password', continuation_context.get('password', ''))}]" | |
| ) | |
| return ( | |
| "Die Registrierung konnte nach dem CAPTCHA-Handoff nicht abgeschlossen werden: " | |
| f"{workflow_result.get('error', 'Unbekannter Fehler')}" | |
| ) | |
| continuation_note = ( | |
| "\n\n[CAPTCHA-HANDOFF FORTSETZEN]\n" | |
| "Der Nutzer hat die Sicherheitsprüfung erfolgreich abgeschlossen. " | |
| "Die aktuelle Browser-Seite und Cookies enthalten die bestätigte Sitzung. " | |
| "Setze die ursprüngliche Aufgabe auf der aktuellen Seite fort; beginne nicht mit einem neuen Browser-Handoff. " | |
| "Prüfe zuerst die aktuelle Seite und führe nur die verbleibenden Schritte aus." | |
| ) | |
| logger.info( | |
| "CAPTCHA-Handoff setzt ursprünglichen Task %s mit bestätigter Sitzung fort.", | |
| source_task_id, | |
| ) | |
| model = source_params.get("model") | |
| if model and model.startswith("openrouter-"): | |
| model = model.replace("openrouter-", "", 1) | |
| return await brain.run( | |
| original_instruction + continuation_note, | |
| page, | |
| model=model, | |
| task=task, | |
| ) | |
| case _: | |
| return f"Unbekannter Task-Typ: '{task.task_type}'" | |
| finally: | |
| if task.save_profile and task.status == TaskStatus.DONE: | |
| try: | |
| new_state = await ctx.storage_state() | |
| encrypted_state = encrypt_state(new_state) | |
| task.steps.append({"type": "profile_exported", "encrypted_state": encrypted_state}) | |
| logger.info("Persistent Profile state encrypted and exported.") | |
| except Exception as e: | |
| logger.error(f"Failed to export profile state: {e}") | |
| await ctx.close() | |
| async def _notify_telegram(bot, task: Task): | |
| icon = "✅" if task.status == TaskStatus.DONE else "❌" | |
| msg = ( | |
| f"{icon} **Task abgeschlossen**\n" | |
| f"ID: `{task.id}`\n" | |
| f"Typ: `{task.task_type}`\n\n" | |
| f"{'Ergebnis:' if task.result else 'Fehler:'}\n" | |
| f"{task.result or task.error}" | |
| ) | |
| # This requires the raw `bot` instance from telegram.ext.Application.bot | |
| await bot.send_message(chat_id=task.chat_id, text=msg, parse_mode="Markdown") | |