# 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" @dataclass 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 @property 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 @staticmethod 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() @property 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")