Luca448's picture
Savestate: Fix Dockerfile and add Browser_Agent for HF build, bump version to 1.8.334
508316a
Raw History Blame Contribute Delete
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"
@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")