Spaces:
Running
Running
Download Browser_Agent/api/server.py from Luca448/APP-Backend: direct link, hf CLI and curl.
- Browser
- Download file 13.2 kB
-
https://huggingface.co/spaces/Luca448/APP-Backend/resolve/main/Browser_Agent/api/server.py
- Command line
-
hf download hf://spaces/Luca448/APP-Backend/Browser_Agent/api/server.py
-
curl -L -o server.py https://huggingface.co/spaces/Luca448/APP-Backend/resolve/main/Browser_Agent/api/server.py
13.2 kB
| # api/server.py | |
| """ | |
| Interne FastAPI für Status-Checks, Task-Enqueuing und ephemeren Medien-Transfer. | |
| Läuft auf Port 8765. Medien verbleiben nur flüchtig im RAM (TTL). | |
| """ | |
| import os | |
| import secrets | |
| from typing import Any | |
| from fastapi import Depends, FastAPI, Header, HTTPException, Response | |
| from fastapi.responses import FileResponse | |
| from pydantic import BaseModel, Field | |
| from config.settings import settings | |
| from job_queue.task_queue import task_queue | |
| from api.models import TaskRequest, TaskResponse | |
| from workflows.ephemeral_media import ephemeral_media_store | |
| from security.crypto import encrypt_bytes | |
| from agent.vision_subagent import VisionCheckError, vision_subagent | |
| app = FastAPI(title="Gravity Browser Agent API", version="0.2.0") | |
| class NewTaskRequest(BaseModel): | |
| task_type: str | |
| params: dict | |
| chat_id: int | None = None | |
| class HandoffCompletionRequest(BaseModel): | |
| revision: int = Field(ge=0) | |
| cookies: list[dict[str, Any]] = Field(default_factory=list) | |
| resume_url: str | |
| submission_confirmed: bool = False | |
| class DecisionRequest(BaseModel): | |
| confirmation_id: str | |
| revision: int = Field(ge=0) | |
| approved: bool | |
| class VisionCheckRequest(BaseModel): | |
| goal: str = Field(min_length=1, max_length=4_000) | |
| images: list[str] = Field(min_length=1, max_length=3) | |
| def require_agent_auth(authorization: str | None = Header(default=None)) -> None: | |
| expected = settings.gravity_token.strip() | |
| if not expected: | |
| if settings.require_api_auth: | |
| raise HTTPException(status_code=503, detail="Agent API authentication is not configured") | |
| return | |
| scheme, _, supplied = (authorization or "").partition(" ") | |
| if scheme.lower() != "bearer" or not secrets.compare_digest(supplied, expected): | |
| raise HTTPException(status_code=401, detail="Invalid agent access token") | |
| async def health(): | |
| ephemeral_media_store.purge_expired() | |
| return {"status": "ok", "queue_size": task_queue.size} | |
| async def ready(): | |
| browser_ready = bool(getattr(__import__("browser.manager", fromlist=["browser_manager"]).browser_manager, "is_ready", False)) | |
| store = task_queue._store | |
| dependencies = { | |
| "worker": { | |
| "ready": task_queue.worker_heartbeat is not None, | |
| "heartbeat_at": task_queue.worker_heartbeat.isoformat() if task_queue.worker_heartbeat else None, | |
| }, | |
| "browser": {"ready": browser_ready}, | |
| "supabase": { | |
| "configured": store.configured, | |
| "ready": store.ready, | |
| "error": store.last_error, | |
| }, | |
| "auth": {"ready": bool(settings.gravity_token) or not settings.require_api_auth}, | |
| "checkpoint_encryption": { | |
| "ready": len((settings.aes_secret_key or settings.gravity_token).strip()) >= 16 | |
| }, | |
| "gmail": {"configured": bool(settings.agent_gmail_user and settings.agent_gmail_password)}, | |
| } | |
| core_ready = ( | |
| dependencies["worker"]["ready"] | |
| and dependencies["auth"]["ready"] | |
| and dependencies["checkpoint_encryption"]["ready"] | |
| ) | |
| return {"status": "ready" if core_ready else "degraded", "dependencies": dependencies} | |
| async def visual_check(req: VisionCheckRequest): | |
| """Analyze ephemeral screenshots without persisting or logging payloads.""" | |
| try: | |
| return await vision_subagent.analyze(req.images, req.goal) | |
| except VisionCheckError as exc: | |
| raise HTTPException( | |
| status_code=503 if exc.retryable else 400, | |
| detail={"code": "VISUAL_CHECK_FAILED", "message": str(exc), "retryable": exc.retryable}, | |
| ) from exc | |
| async def close_vision_client(): | |
| await vision_subagent.close() | |
| async def create_task( | |
| req: NewTaskRequest, | |
| idempotency_key: str | None = Header(default=None, alias="Idempotency-Key"), | |
| ): | |
| # One-release legacy adapter: old Flutter builds submit captcha_resume as a | |
| # new task. Resolve it onto the original waiting task instead. | |
| if req.task_type == "captcha_resume": | |
| source_task_id = str(req.params.get("source_task_id") or "") | |
| source_task = task_queue._tasks.get(source_task_id) | |
| if source_task and source_task.handoff: | |
| try: | |
| await task_queue.resume_handoff( | |
| source_task_id, | |
| str(source_task.handoff["id"]), | |
| revision=source_task.revision, | |
| cookies=list(req.params.get("exported_cookies") or []), | |
| resume_url=str(req.params.get("resume_url") or source_task.handoff.get("url") or ""), | |
| submission_confirmed=req.params.get("submission_confirmed") is True, | |
| ) | |
| return {"task_id": source_task_id, "legacy_resume": True} | |
| except (KeyError, ValueError) as exc: | |
| raise HTTPException(status_code=409, detail=str(exc)) from exc | |
| try: | |
| task_id = await task_queue.enqueue( | |
| req.task_type, | |
| req.params, | |
| req.chat_id, | |
| idempotency_key=idempotency_key, | |
| ) | |
| except ValueError as exc: | |
| if str(exc) == "idempotency_conflict": | |
| raise HTTPException(status_code=409, detail="Idempotency key payload conflict") from exc | |
| raise | |
| return {"task_id": task_id} | |
| async def get_task_status(task_id: str): | |
| status = await task_queue.get_status(task_id) | |
| if "id" not in status: | |
| raise HTTPException(status_code=404, detail=status.get("error", "Not found")) | |
| return status | |
| async def cancel_task(task_id: str): | |
| cancelled = await task_queue.cancel(task_id) | |
| if not cancelled: | |
| raise HTTPException(status_code=404, detail="Task nicht gefunden") | |
| return {"id": task_id, "status": "cancelled"} | |
| async def create_task_v1( | |
| req: NewTaskRequest, | |
| idempotency_key: str = Header(alias="Idempotency-Key"), | |
| ): | |
| if not idempotency_key.strip(): | |
| raise HTTPException(status_code=400, detail="Idempotency-Key is required") | |
| try: | |
| task_id = await task_queue.enqueue( | |
| req.task_type, | |
| req.params, | |
| req.chat_id, | |
| idempotency_key=idempotency_key, | |
| ) | |
| except ValueError as exc: | |
| if str(exc) == "idempotency_conflict": | |
| raise HTTPException(status_code=409, detail="Idempotency key payload conflict") from exc | |
| raise | |
| return await task_queue.get_status(task_id, typed=True) | |
| async def get_task_v1(task_id: str): | |
| status = await task_queue.get_status(task_id, typed=True) | |
| if "id" not in status: | |
| raise HTTPException(status_code=404, detail=status.get("error", "Not found")) | |
| return status | |
| async def get_task_events(task_id: str, after: int = 0): | |
| task = task_queue._tasks.get(task_id) | |
| if not task: | |
| raise HTTPException(status_code=404, detail="Task not found") | |
| return { | |
| "task_id": task_id, | |
| "revision": task.revision, | |
| "events": [event for event in task.events if int(event.get("sequence", 0)) > after], | |
| } | |
| async def get_handoff_payload(task_id: str, handoff_id: str): | |
| payload = await task_queue.get_handoff_payload(task_id, handoff_id) | |
| if not payload: | |
| raise HTTPException(status_code=410, detail="Handoff expired or unavailable") | |
| return payload | |
| async def complete_handoff(task_id: str, handoff_id: str, req: HandoffCompletionRequest): | |
| try: | |
| task = await task_queue.resume_handoff( | |
| task_id, | |
| handoff_id, | |
| revision=req.revision, | |
| cookies=req.cookies, | |
| resume_url=req.resume_url, | |
| submission_confirmed=req.submission_confirmed, | |
| ) | |
| except KeyError as exc: | |
| raise HTTPException(status_code=404, detail=str(exc)) from exc | |
| except ValueError as exc: | |
| code = 410 if str(exc) == "handoff_expired" else 409 | |
| raise HTTPException(status_code=code, detail=str(exc)) from exc | |
| return await task_queue.get_status(task.id, typed=True) | |
| async def decide_task(task_id: str, req: DecisionRequest): | |
| try: | |
| task = await task_queue.resolve_confirmation( | |
| task_id, | |
| confirmation_id=req.confirmation_id, | |
| revision=req.revision, | |
| approved=req.approved, | |
| ) | |
| except KeyError as exc: | |
| raise HTTPException(status_code=404, detail=str(exc)) from exc | |
| except ValueError as exc: | |
| code = 410 if str(exc) == "confirmation_expired" else 409 | |
| raise HTTPException(status_code=code, detail=str(exc)) from exc | |
| return await task_queue.get_status(task.id, typed=True) | |
| async def cancel_task_v1(task_id: str): | |
| if not await task_queue.cancel(task_id): | |
| raise HTTPException(status_code=404, detail="Task not found") | |
| return await task_queue.get_status(task_id, typed=True) | |
| # ── Ephemere Medien Endpunkte (RAM-only, keine Festplattenspeicherung) ────────── | |
| async def get_ephemeral_media(media_id: str): | |
| """ | |
| Liefert das temporär im RAM gehaltene Bild direkt aus. | |
| Headers verhindern Caching auf Zwischenservern. | |
| """ | |
| item = ephemeral_media_store.get(media_id) | |
| if not item: | |
| raise HTTPException(status_code=404, detail="Medienobjekt abgelaufen oder gelöscht") | |
| headers = { | |
| "Cache-Control": "no-store, no-cache, must-revalidate, max-age=0", | |
| "Pragma": "no-cache", | |
| "Content-Disposition": f'inline; filename="{item.filename}"', | |
| "X-Content-Type-Options": "nosniff", | |
| "Referrer-Policy": "no-referrer", | |
| "Cross-Origin-Resource-Policy": "same-origin", | |
| "Permissions-Policy": "interest-cohort=()", | |
| } | |
| return Response(content=item.media_bytes, media_type=item.media_type, headers=headers) | |
| async def get_encrypted_media(media_id: str): | |
| """ | |
| Liefert die Bild-Payload AES-verschlüsselt für den sicheren Mobile-Tresor. | |
| """ | |
| item = ephemeral_media_store.get(media_id) | |
| if not item: | |
| raise HTTPException(status_code=404, detail="Medienobjekt abgelaufen oder gelöscht") | |
| encrypted_payload = encrypt_bytes(item.media_bytes) | |
| headers = { | |
| "Cache-Control": "no-store, no-cache, must-revalidate", | |
| "Content-Disposition": f'attachment; filename="{item.filename}.enc"', | |
| "X-Content-Type-Options": "nosniff", | |
| "Referrer-Policy": "no-referrer", | |
| "Cross-Origin-Resource-Policy": "same-origin", | |
| "Permissions-Policy": "interest-cohort=()", | |
| } | |
| return Response(content=encrypted_payload, media_type="application/octet-stream", headers=headers) | |
| async def delete_ephemeral_media(media_id: str): | |
| """ | |
| Vernichtet das Bild sofort und unwiderruflich aus dem Server-Arbeitsspeicher. | |
| """ | |
| deleted = ephemeral_media_store.delete(media_id) | |
| return {"deleted": deleted} | |
| async def get_download(filename: str): | |
| downloads_root = settings.downloads_dir.resolve() | |
| file_path = (downloads_root / filename).resolve() | |
| if not file_path.is_relative_to(downloads_root): | |
| raise HTTPException(status_code=403, detail="Forbidden") | |
| if not file_path.exists() or not file_path.is_file(): | |
| raise HTTPException(status_code=404, detail="File not found") | |
| return FileResponse(path=file_path, filename=filename) | |
| async def run_task(request: TaskRequest): | |
| task_desc = request.task | |
| if request.predefined_task: | |
| from workflows.tasks import get_predefined_task | |
| predef = get_predefined_task(request.predefined_task) | |
| if predef: | |
| task_desc = predef + "\n" + task_desc | |
| if not task_desc: | |
| raise HTTPException(status_code=400, detail="Kein Task angegeben.") | |
| task_id = await task_queue.enqueue("agent_task", {"task": task_desc}) | |
| return TaskResponse(status="pending", result=f"Task eingereiht mit ID: {task_id}") | |