| 1 | from __future__ import annotations |
| 2 | |
| 3 | import asyncio |
| 4 | import base64 |
| 5 | import contextlib |
| 6 | import time |
| 7 | from typing import Any, ClassVar |
| 8 | |
| 9 | from agent import AgentContext |
| 10 | from helpers.ws import WsHandler |
| 11 | from helpers.ws_manager import WsResult |
| 12 | from plugins._browser.helpers.config import ( |
| 13 | DEFAULT_BROWSER_TAB_SCOPE, |
| 14 | TAB_SCOPE_KEY, |
| 15 | get_browser_config, |
| 16 | ) |
| 17 | from plugins._browser.helpers.runtime import ( |
| 18 | get_runtime, |
| 19 | has_restorable_browser_tabs, |
| 20 | list_runtime_sessions, |
| 21 | ) |
| 22 | |
| 23 | |
| 24 | FRAME_READ_TIMEOUT_SECONDS = 0.5 |
| 25 | FRAME_RETRY_DELAY_SECONDS = 0.5 |
| 26 | FRAME_STATE_REFRESH_SECONDS = 0.75 |
| 27 | SNAPSHOT_STATE_POLL_SECONDS = 0.75 |
| 28 | SCREENCAST_STREAM_QUALITY = 80 |
| 29 | SCREENSHOT_QUALITY = 92 |
| 30 | VIEWER_TRANSPORT_SCREENCAST = "screencast" |
| 31 | VIEWER_TRANSPORT_SNAPSHOT = "snapshot" |
| 32 | VIEWER_TRANSPORT_INTERACTIVE = "interactive" |
| 33 | VIEWER_TRANSPORTS = { |
| 34 | VIEWER_TRANSPORT_INTERACTIVE, |
| 35 | VIEWER_TRANSPORT_SCREENCAST, |
| 36 | VIEWER_TRANSPORT_SNAPSHOT, |
| 37 | } |
| 38 | |
| 39 | |
| 40 | class WsBrowser(WsHandler): |
| 41 | _streams: ClassVar[dict[tuple[str, str], asyncio.Task[None]]] = {} |
| 42 | |
| 43 | async def on_disconnect(self, sid: str) -> None: |
| 44 | for key in [key for key in self._streams if key[0] == sid]: |
| 45 | task = self._streams.pop(key) |
| 46 | task.cancel() |
| 47 | |
| 48 | async def process( |
| 49 | self, |
| 50 | event: str, |
| 51 | data: dict[str, Any], |
| 52 | sid: str, |
| 53 | ) -> dict[str, Any] | WsResult | None: |
| 54 | if not event.startswith("browser_"): |
| 55 | return None |
| 56 | |
| 57 | if event == "browser_viewer_subscribe": |
| 58 | return await self._subscribe(data, sid) |
| 59 | if event == "browser_viewer_unsubscribe": |
| 60 | return self._unsubscribe(data, sid) |
| 61 | if event == "browser_viewer_snapshot": |
| 62 | return await self._snapshot(data) |
| 63 | if event == "browser_viewer_sessions": |
| 64 | return await self._sessions(data) |
| 65 | if event == "browser_viewer_command": |
| 66 | return await self._command(data, sid) |
| 67 | if event == "browser_viewer_input": |
| 68 | return await self._input(data, sid) |
| 69 | if event == "browser_viewer_annotation": |
| 70 | return await self._annotation(data, sid) |
| 71 | |
| 72 | return WsResult.error( |
| 73 | code="UNKNOWN_BROWSER_EVENT", |
| 74 | message=f"Unknown browser event: {event}", |
| 75 | correlation_id=data.get("correlationId"), |
| 76 | ) |
| 77 | |
| 78 | async def _subscribe(self, data: dict[str, Any], sid: str) -> dict[str, Any] | WsResult: |
| 79 | context_id = self._context_id(data) |
| 80 | if not context_id: |
| 81 | return self._error("MISSING_CONTEXT", "context_id is required", data) |
| 82 | if not AgentContext.get(context_id): |
| 83 | return self._error("CONTEXT_NOT_FOUND", f"Context '{context_id}' was not found", data) |
| 84 | |
| 85 | create_browser = self._bool(data.get("create_browser", data.get("createBrowser"))) |
| 86 | runtime = await get_runtime(context_id, create=create_browser) |
| 87 | if not runtime and not create_browser and has_restorable_browser_tabs(context_id): |
| 88 | runtime = await get_runtime(context_id) |
| 89 | listing = {"browsers": [], "last_interacted_browser_id": None} |
| 90 | browsers: list[dict[str, Any]] = [] |
| 91 | if runtime: |
| 92 | listing = await runtime.call("list") |
| 93 | browsers = listing.get("browsers") or [] |
| 94 | if runtime and not browsers and create_browser: |
| 95 | opened = await runtime.call("open", "") |
| 96 | listing = await runtime.call("list") |
| 97 | browsers = listing.get("browsers") or [] |
| 98 | if opened.get("id"): |
| 99 | listing["last_interacted_browser_id"] = opened.get("id") |
| 100 | active_id = self._active_browser_id(listing, data.get("browser_id")) |
| 101 | requested_transport = self._viewer_transport(data) |
| 102 | viewer_transport, interactive_view = await self._effective_viewer( |
| 103 | runtime, |
| 104 | active_id, |
| 105 | data, |
| 106 | ) |
| 107 | initial_viewport = self._viewport_from_data(data) |
| 108 | if ( |
| 109 | runtime |
| 110 | and active_id |
| 111 | and initial_viewport |
| 112 | and viewer_transport != VIEWER_TRANSPORT_INTERACTIVE |
| 113 | ): |
| 114 | await runtime.call( |
| 115 | "set_viewport", |
| 116 | active_id, |
| 117 | initial_viewport["width"], |
| 118 | initial_viewport["height"], |
| 119 | ) |
| 120 | listing = await runtime.call("list") |
| 121 | browsers = listing.get("browsers") or [] |
| 122 | |
| 123 | stream_key = (sid, context_id) |
| 124 | existing = self._streams.pop(stream_key, None) |
| 125 | if existing: |
| 126 | existing.cancel() |
| 127 | viewer_id = str(data.get("viewer_id") or "") |
| 128 | binary_frames = self._bool(data.get("binary_frames", data.get("binaryFrames"))) |
| 129 | slim_frames = self._bool(data.get("slim_frames", data.get("slimFrames", binary_frames))) |
| 130 | capture_scale = self._capture_scale_from_data(data) |
| 131 | snapshot = None |
| 132 | if runtime: |
| 133 | if viewer_transport == VIEWER_TRANSPORT_SCREENCAST: |
| 134 | stream_task = self._stream_frames( |
| 135 | sid, |
| 136 | context_id, |
| 137 | active_id, |
| 138 | viewer_id, |
| 139 | binary_frames=binary_frames, |
| 140 | slim_frames=slim_frames, |
| 141 | capture_scale=capture_scale, |
| 142 | ) |
| 143 | else: |
| 144 | stream_task = self._stream_state( |
| 145 | sid, |
| 146 | context_id, |
| 147 | active_id, |
| 148 | viewer_id, |
| 149 | viewer_transport=viewer_transport, |
| 150 | ) |
| 151 | self._streams[stream_key] = asyncio.create_task(stream_task) |
| 152 | if viewer_transport != VIEWER_TRANSPORT_INTERACTIVE: |
| 153 | snapshot = await self._snapshot_for_browser(runtime, active_id) |
| 154 | |
| 155 | browsers, all_browsers, tab_scope = await self._tabs_for_scope(context_id, browsers) |
| 156 | |
| 157 | return { |
| 158 | "context_id": context_id, |
| 159 | "active_browser_context_id": context_id, |
| 160 | "active_browser_id": active_id, |
| 161 | "snapshot": snapshot, |
| 162 | "browsers": browsers, |
| 163 | "all_browsers": all_browsers, |
| 164 | "tab_scope": tab_scope, |
| 165 | "viewer_id": viewer_id, |
| 166 | "viewer_transport": viewer_transport, |
| 167 | "interactive_view": interactive_view, |
| 168 | "viewer_fallback_reason": ( |
| 169 | str(interactive_view.get("error") or "") |
| 170 | if requested_transport == VIEWER_TRANSPORT_INTERACTIVE |
| 171 | and interactive_view |
| 172 | and not interactive_view.get("available") |
| 173 | else "" |
| 174 | ), |
| 175 | "binary_frames": binary_frames, |
| 176 | "slim_frames": slim_frames, |
| 177 | } |
| 178 | |
| 179 | def _unsubscribe(self, data: dict[str, Any], sid: str) -> dict[str, Any] | WsResult: |
| 180 | context_id = self._context_id(data) |
| 181 | if not context_id: |
| 182 | return self._error("MISSING_CONTEXT", "context_id is required", data) |
| 183 | task = self._streams.pop((sid, context_id), None) |
| 184 | if task: |
| 185 | task.cancel() |
| 186 | return {"context_id": context_id, "unsubscribed": True} |
| 187 | |
| 188 | async def _sessions(self, data: dict[str, Any]) -> dict[str, Any]: |
| 189 | context_id = self._context_id(data) |
| 190 | tab_scope = self._tab_scope() |
| 191 | if tab_scope == "shared": |
| 192 | return { |
| 193 | "context_id": context_id, |
| 194 | "browsers": await self._all_browser_tabs(), |
| 195 | "all_browsers": True, |
| 196 | "tab_scope": tab_scope, |
| 197 | } |
| 198 | |
| 199 | runtime = await get_runtime(context_id, create=False) if context_id else None |
| 200 | listing = await runtime.call("list") if runtime else {} |
| 201 | return { |
| 202 | "context_id": context_id, |
| 203 | "browsers": listing.get("browsers") or [], |
| 204 | "all_browsers": False, |
| 205 | "tab_scope": tab_scope, |
| 206 | } |
| 207 | |
| 208 | async def _snapshot(self, data: dict[str, Any]) -> dict[str, Any] | WsResult: |
| 209 | context_id = self._context_id(data) |
| 210 | if not context_id: |
| 211 | return self._error("MISSING_CONTEXT", "context_id is required", data) |
| 212 | if not AgentContext.get(context_id): |
| 213 | return self._error("CONTEXT_NOT_FOUND", f"Context '{context_id}' was not found", data) |
| 214 | |
| 215 | runtime = await get_runtime(context_id, create=False) |
| 216 | if not runtime: |
| 217 | browsers, all_browsers, tab_scope = await self._tabs_for_scope(context_id, []) |
| 218 | return { |
| 219 | "context_id": context_id, |
| 220 | "active_browser_context_id": context_id, |
| 221 | "active_browser_id": None, |
| 222 | "snapshot": None, |
| 223 | "browsers": browsers, |
| 224 | "all_browsers": all_browsers, |
| 225 | "tab_scope": tab_scope, |
| 226 | } |
| 227 | |
| 228 | listing = await runtime.call("list") |
| 229 | browsers = listing.get("browsers") or [] |
| 230 | active_id = self._active_browser_id(listing, data.get("browser_id")) |
| 231 | snapshot = None |
| 232 | if active_id: |
| 233 | try: |
| 234 | quality = int(data.get("quality") or SCREENSHOT_QUALITY) |
| 235 | except (TypeError, ValueError): |
| 236 | quality = SCREENSHOT_QUALITY |
| 237 | with contextlib.suppress(Exception): |
| 238 | snapshot = await runtime.call( |
| 239 | "screenshot", |
| 240 | active_id, |
| 241 | quality=quality, |
| 242 | ) |
| 243 | |
| 244 | browsers, all_browsers, tab_scope = await self._tabs_for_scope(context_id, browsers) |
| 245 | |
| 246 | return { |
| 247 | "context_id": context_id, |
| 248 | "active_browser_context_id": context_id, |
| 249 | "active_browser_id": active_id, |
| 250 | "snapshot": snapshot, |
| 251 | "browsers": browsers, |
| 252 | "all_browsers": all_browsers, |
| 253 | "tab_scope": tab_scope, |
| 254 | } |
| 255 | |
| 256 | async def _command(self, data: dict[str, Any], sid: str) -> dict[str, Any] | WsResult: |
| 257 | context_id = self._context_id(data) |
| 258 | if not context_id: |
| 259 | return self._error("MISSING_CONTEXT", "context_id is required", data) |
| 260 | runtime = await get_runtime(context_id) |
| 261 | command = str(data.get("command") or "").strip().lower().replace("-", "_") |
| 262 | browser_id = data.get("browser_id") |
| 263 | viewer_id = str(data.get("viewer_id") or "") |
| 264 | |
| 265 | try: |
| 266 | if command == "open": |
| 267 | result = await runtime.call("open", data.get("url") or "") |
| 268 | elif command == "navigate": |
| 269 | result = await runtime.call( |
| 270 | "navigate", |
| 271 | browser_id, |
| 272 | data.get("url") or "", |
| 273 | wait_until="commit", |
| 274 | ) |
| 275 | elif command == "back": |
| 276 | result = await runtime.call("back", browser_id, wait_until="commit") |
| 277 | elif command == "forward": |
| 278 | result = await runtime.call("forward", browser_id, wait_until="commit") |
| 279 | elif command == "reload": |
| 280 | result = await runtime.call("reload", browser_id, wait_until="commit") |
| 281 | elif command == "close": |
| 282 | result = await runtime.call("close_browser", browser_id) |
| 283 | elif command == "list": |
| 284 | result = await runtime.call("list") |
| 285 | else: |
| 286 | return self._error("UNKNOWN_COMMAND", f"Unknown browser command: {command}", data) |
| 287 | except Exception as exc: |
| 288 | return self._error("COMMAND_FAILED", str(exc), data) |
| 289 | |
| 290 | listing = await runtime.call("list") |
| 291 | last_interacted_browser_id = listing.get("last_interacted_browser_id") |
| 292 | active_id = self._active_browser_id( |
| 293 | listing, |
| 294 | self._result_browser_id(result) or browser_id, |
| 295 | ) |
| 296 | viewer_transport, interactive_view = await self._effective_viewer( |
| 297 | runtime, |
| 298 | active_id, |
| 299 | data, |
| 300 | ) |
| 301 | snapshot = ( |
| 302 | None |
| 303 | if viewer_transport == VIEWER_TRANSPORT_INTERACTIVE |
| 304 | else await self._snapshot_for_result(runtime, result) |
| 305 | ) |
| 306 | browsers, all_browsers, tab_scope = await self._tabs_for_scope( |
| 307 | context_id, |
| 308 | listing.get("browsers") or [], |
| 309 | ) |
| 310 | await self.emit_to( |
| 311 | sid, |
| 312 | "browser_viewer_state", |
| 313 | { |
| 314 | "context_id": context_id, |
| 315 | "active_browser_context_id": context_id, |
| 316 | "viewer_id": viewer_id, |
| 317 | "command": command, |
| 318 | "browser_id": browser_id, |
| 319 | "result": result, |
| 320 | "snapshot": snapshot, |
| 321 | "browsers": browsers, |
| 322 | "all_browsers": all_browsers, |
| 323 | "tab_scope": tab_scope, |
| 324 | "last_interacted_browser_id": last_interacted_browser_id, |
| 325 | "viewer_transport": viewer_transport, |
| 326 | "interactive_view": interactive_view, |
| 327 | }, |
| 328 | correlation_id=data.get("correlationId"), |
| 329 | ) |
| 330 | return { |
| 331 | "result": result, |
| 332 | "snapshot": snapshot, |
| 333 | "browsers": browsers, |
| 334 | "all_browsers": all_browsers, |
| 335 | "tab_scope": tab_scope, |
| 336 | "active_browser_context_id": context_id, |
| 337 | "last_interacted_browser_id": last_interacted_browser_id, |
| 338 | "command": command, |
| 339 | "browser_id": browser_id, |
| 340 | "viewer_id": viewer_id, |
| 341 | "viewer_transport": viewer_transport, |
| 342 | "interactive_view": interactive_view, |
| 343 | } |
| 344 | |
| 345 | async def _input(self, data: dict[str, Any], sid: str) -> dict[str, Any] | WsResult: |
| 346 | context_id = self._context_id(data) |
| 347 | if not context_id: |
| 348 | return self._error("MISSING_CONTEXT", "context_id is required", data) |
| 349 | runtime = await get_runtime(context_id, create=False) |
| 350 | if not runtime: |
| 351 | return self._error("NO_BROWSER_RUNTIME", "No browser runtime exists for this context", data) |
| 352 | |
| 353 | input_type = str(data.get("input_type") or "").strip().lower() |
| 354 | browser_id = data.get("browser_id") |
| 355 | try: |
| 356 | if input_type == "mouse": |
| 357 | result = await runtime.call( |
| 358 | "mouse", |
| 359 | browser_id, |
| 360 | data.get("event_type") or "click", |
| 361 | float(data.get("x") or 0), |
| 362 | float(data.get("y") or 0), |
| 363 | data.get("button") or "left", |
| 364 | ) |
| 365 | elif input_type == "keyboard": |
| 366 | result = await runtime.call( |
| 367 | "keyboard", |
| 368 | browser_id, |
| 369 | key=str(data.get("key") or ""), |
| 370 | text=str(data.get("text") or ""), |
| 371 | ) |
| 372 | elif input_type == "clipboard": |
| 373 | result = await runtime.call( |
| 374 | "clipboard", |
| 375 | browser_id, |
| 376 | action=str(data.get("action") or ""), |
| 377 | text=str(data.get("text") or ""), |
| 378 | ) |
| 379 | elif input_type == "viewport": |
| 380 | viewer_transport = self._viewer_transport(data) |
| 381 | result = await runtime.call( |
| 382 | "set_viewport", |
| 383 | browser_id, |
| 384 | int(data.get("width") or 0), |
| 385 | int(data.get("height") or 0), |
| 386 | restart_screencast=bool(data.get("restart_stream")), |
| 387 | resize_interactive=viewer_transport == VIEWER_TRANSPORT_INTERACTIVE, |
| 388 | include_state=viewer_transport != VIEWER_TRANSPORT_INTERACTIVE, |
| 389 | ) |
| 390 | elif input_type == "wheel": |
| 391 | result = await runtime.call( |
| 392 | "wheel", |
| 393 | browser_id, |
| 394 | float(data.get("x") or 0), |
| 395 | float(data.get("y") or 0), |
| 396 | float(data.get("delta_x") or 0), |
| 397 | float(data.get("delta_y") or 0), |
| 398 | ) |
| 399 | else: |
| 400 | return self._error("UNKNOWN_INPUT", f"Unknown browser input: {input_type}", data) |
| 401 | except Exception as exc: |
| 402 | return self._error("INPUT_FAILED", str(exc), data) |
| 403 | |
| 404 | if input_type == "clipboard": |
| 405 | response = { |
| 406 | "state": result.get("state") if isinstance(result, dict) else result, |
| 407 | "snapshot": None, |
| 408 | } |
| 409 | if isinstance(result, dict): |
| 410 | response["clipboard"] = result.get("clipboard") |
| 411 | return response |
| 412 | |
| 413 | return { |
| 414 | "state": result, |
| 415 | "snapshot": await self._snapshot_for_result(runtime, result) |
| 416 | if input_type == "mouse" |
| 417 | else None, |
| 418 | } |
| 419 | |
| 420 | async def _annotation(self, data: dict[str, Any], sid: str) -> dict[str, Any] | WsResult: |
| 421 | context_id = self._context_id(data) |
| 422 | if not context_id: |
| 423 | return self._error("MISSING_CONTEXT", "context_id is required", data) |
| 424 | runtime = await get_runtime(context_id, create=False) |
| 425 | if not runtime: |
| 426 | return self._error("NO_BROWSER_RUNTIME", "No browser runtime exists for this context", data) |
| 427 | |
| 428 | browser_id = data.get("browser_id") |
| 429 | viewer_id = str(data.get("viewer_id") or "") |
| 430 | payload = data.get("payload") if isinstance(data.get("payload"), dict) else {} |
| 431 | try: |
| 432 | annotation = await runtime.call("annotation_target", browser_id, payload) |
| 433 | except Exception as exc: |
| 434 | return self._error("ANNOTATION_FAILED", str(exc), data) |
| 435 | |
| 436 | return { |
| 437 | "annotation": annotation, |
| 438 | "context_id": context_id, |
| 439 | "browser_id": browser_id, |
| 440 | "viewer_id": viewer_id, |
| 441 | } |
| 442 | |
| 443 | async def _snapshot_for_result( |
| 444 | self, |
| 445 | runtime: Any, |
| 446 | result: dict[str, Any] | None, |
| 447 | ) -> dict[str, Any] | None: |
| 448 | if not isinstance(result, dict): |
| 449 | return None |
| 450 | state = result.get("state") if isinstance(result.get("state"), dict) else result |
| 451 | browser_id = state.get("id") if isinstance(state, dict) else result.get("id") |
| 452 | if not browser_id: |
| 453 | return None |
| 454 | with contextlib.suppress(Exception): |
| 455 | return await runtime.call("screenshot", browser_id, quality=SCREENSHOT_QUALITY) |
| 456 | return None |
| 457 | |
| 458 | async def _snapshot_for_browser( |
| 459 | self, |
| 460 | runtime: Any, |
| 461 | browser_id: int | str | None, |
| 462 | ) -> dict[str, Any] | None: |
| 463 | if not browser_id: |
| 464 | return None |
| 465 | with contextlib.suppress(Exception): |
| 466 | return await runtime.call("screenshot", browser_id, quality=SCREENSHOT_QUALITY) |
| 467 | return None |
| 468 | |
| 469 | @staticmethod |
| 470 | def _tab_scope() -> str: |
| 471 | scope = str( |
| 472 | (get_browser_config() or {}).get(TAB_SCOPE_KEY, DEFAULT_BROWSER_TAB_SCOPE) |
| 473 | or DEFAULT_BROWSER_TAB_SCOPE |
| 474 | ).strip().lower().replace("-", "_") |
| 475 | return "shared" if scope == "shared" else DEFAULT_BROWSER_TAB_SCOPE |
| 476 | |
| 477 | async def _tabs_for_scope( |
| 478 | self, |
| 479 | context_id: str, |
| 480 | browsers: list[dict[str, Any]] | None, |
| 481 | ) -> tuple[list[dict[str, Any]], bool, str]: |
| 482 | tab_scope = self._tab_scope() |
| 483 | if tab_scope == "shared": |
| 484 | return await self._all_browser_tabs(), True, tab_scope |
| 485 | return browsers or [], False, tab_scope |
| 486 | |
| 487 | async def _all_browser_tabs(self) -> list[dict[str, Any]]: |
| 488 | browsers: list[dict[str, Any]] = [] |
| 489 | for session in await list_runtime_sessions(): |
| 490 | context_id = str(session.get("context_id") or "") |
| 491 | for browser in session.get("browsers") or []: |
| 492 | entry = dict(browser or {}) |
| 493 | entry.setdefault("context_id", context_id) |
| 494 | browsers.append(entry) |
| 495 | return browsers |
| 496 | |
| 497 | async def _stream_frames( |
| 498 | self, |
| 499 | sid: str, |
| 500 | context_id: str, |
| 501 | browser_id: int | str | None, |
| 502 | viewer_id: str = "", |
| 503 | *, |
| 504 | binary_frames: bool = False, |
| 505 | slim_frames: bool = False, |
| 506 | capture_scale: float = 1.0, |
| 507 | ) -> None: |
| 508 | runtime = None |
| 509 | stream_id = None |
| 510 | while True: |
| 511 | try: |
| 512 | runtime = await get_runtime(context_id, create=False) |
| 513 | if not runtime: |
| 514 | await self._emit_empty_frame( |
| 515 | sid, |
| 516 | context_id, |
| 517 | viewer_id=viewer_id, |
| 518 | frame_source=VIEWER_TRANSPORT_SCREENCAST, |
| 519 | ) |
| 520 | await asyncio.sleep(FRAME_RETRY_DELAY_SECONDS) |
| 521 | continue |
| 522 | |
| 523 | listing = await runtime.call("list") |
| 524 | browsers = listing.get("browsers") or [] |
| 525 | active_id = self._active_browser_id(listing, browser_id) |
| 526 | if not active_id: |
| 527 | await self._emit_viewer_state( |
| 528 | sid, |
| 529 | context_id, |
| 530 | active_id, |
| 531 | browsers=browsers, |
| 532 | viewer_id=viewer_id, |
| 533 | state=None, |
| 534 | viewer_transport=VIEWER_TRANSPORT_SCREENCAST, |
| 535 | ) |
| 536 | await asyncio.sleep(FRAME_RETRY_DELAY_SECONDS) |
| 537 | continue |
| 538 | |
| 539 | screencast = await runtime.call( |
| 540 | "start_screencast", |
| 541 | active_id, |
| 542 | quality=SCREENCAST_STREAM_QUALITY, |
| 543 | every_nth_frame=1, |
| 544 | capture_scale=capture_scale, |
| 545 | ) |
| 546 | stream_id = screencast["stream_id"] |
| 547 | active_id = screencast["browser_id"] |
| 548 | state = screencast.get("state") |
| 549 | await self._emit_viewer_state( |
| 550 | sid, |
| 551 | context_id, |
| 552 | active_id, |
| 553 | browsers=browsers, |
| 554 | viewer_id=viewer_id, |
| 555 | state=state, |
| 556 | viewer_transport=VIEWER_TRANSPORT_SCREENCAST, |
| 557 | ) |
| 558 | |
| 559 | last_state_refresh = 0.0 |
| 560 | last_state_signature = self._state_signature(active_id, browsers) |
| 561 | frame_sequence = 0 |
| 562 | server_loop = asyncio.get_running_loop() |
| 563 | stop_event = asyncio.Event() |
| 564 | |
| 565 | async def emit_frame(frame: dict[str, Any]) -> None: |
| 566 | nonlocal frame_sequence |
| 567 | try: |
| 568 | frame_sequence += 1 |
| 569 | payload = self._frame_payload( |
| 570 | frame, |
| 571 | context_id=context_id, |
| 572 | viewer_id=viewer_id, |
| 573 | browser_id=active_id, |
| 574 | sequence=frame_sequence, |
| 575 | binary_frames=binary_frames, |
| 576 | ) |
| 577 | if not slim_frames: |
| 578 | payload["browsers"] = browsers |
| 579 | payload["state"] = state |
| 580 | await self._emit_to_connected_viewer(sid, "browser_viewer_frame", payload) |
| 581 | except BaseException: |
| 582 | stop_event.set() |
| 583 | raise |
| 584 | |
| 585 | def frame_consumer(frame: dict[str, Any]): |
| 586 | return asyncio.run_coroutine_threadsafe(emit_frame(frame), server_loop) |
| 587 | |
| 588 | def stop_consumer() -> None: |
| 589 | server_loop.call_soon_threadsafe(stop_event.set) |
| 590 | |
| 591 | await runtime.call("attach_screencast_consumer", stream_id, frame_consumer, stop_consumer) |
| 592 | |
| 593 | while True: |
| 594 | if stop_event.is_set(): |
| 595 | break |
| 596 | now = time.monotonic() |
| 597 | if now - last_state_refresh >= FRAME_STATE_REFRESH_SECONDS: |
| 598 | listing = await runtime.call("list") |
| 599 | browsers = listing.get("browsers") or [] |
| 600 | browser_ids = {str(browser.get("id")) for browser in browsers} |
| 601 | if str(active_id) not in browser_ids: |
| 602 | break |
| 603 | state = self._state_for_browser(browsers, active_id, state) |
| 604 | state_signature = self._state_signature(active_id, browsers) |
| 605 | if state_signature != last_state_signature: |
| 606 | await self._emit_viewer_state( |
| 607 | sid, |
| 608 | context_id, |
| 609 | active_id, |
| 610 | browsers=browsers, |
| 611 | viewer_id=viewer_id, |
| 612 | state=state, |
| 613 | viewer_transport=VIEWER_TRANSPORT_SCREENCAST, |
| 614 | ) |
| 615 | last_state_signature = state_signature |
| 616 | last_state_refresh = now |
| 617 | |
| 618 | try: |
| 619 | await asyncio.wait_for(stop_event.wait(), timeout=FRAME_READ_TIMEOUT_SECONDS) |
| 620 | break |
| 621 | except TimeoutError: |
| 622 | continue |
| 623 | except asyncio.CancelledError: |
| 624 | raise |
| 625 | except Exception: |
| 626 | await asyncio.sleep(FRAME_RETRY_DELAY_SECONDS) |
| 627 | finally: |
| 628 | if runtime and stream_id: |
| 629 | with contextlib.suppress(Exception): |
| 630 | await runtime.call("stop_screencast", stream_id) |
| 631 | stream_id = None |
| 632 | |
| 633 | async def _stream_state( |
| 634 | self, |
| 635 | sid: str, |
| 636 | context_id: str, |
| 637 | browser_id: int | str | None, |
| 638 | viewer_id: str = "", |
| 639 | *, |
| 640 | viewer_transport: str = VIEWER_TRANSPORT_SNAPSHOT, |
| 641 | ) -> None: |
| 642 | last_signature = None |
| 643 | while True: |
| 644 | try: |
| 645 | runtime = await get_runtime(context_id, create=False) |
| 646 | if not runtime: |
| 647 | signature = (None, ()) |
| 648 | if signature != last_signature: |
| 649 | await self._emit_empty_frame( |
| 650 | sid, |
| 651 | context_id, |
| 652 | viewer_id=viewer_id, |
| 653 | frame_source=viewer_transport, |
| 654 | ) |
| 655 | last_signature = signature |
| 656 | await asyncio.sleep(FRAME_RETRY_DELAY_SECONDS) |
| 657 | continue |
| 658 | |
| 659 | listing = await runtime.call("list") |
| 660 | browsers = listing.get("browsers") or [] |
| 661 | active_id = self._active_browser_id(listing, browser_id) |
| 662 | state = self._state_for_browser(browsers, active_id, None) if active_id else None |
| 663 | signature = ( |
| 664 | str(active_id or ""), |
| 665 | tuple( |
| 666 | ( |
| 667 | str(browser.get("context_id") or context_id), |
| 668 | str(browser.get("id") or ""), |
| 669 | str(browser.get("currentUrl") or ""), |
| 670 | str(browser.get("title") or ""), |
| 671 | bool(browser.get("loading")), |
| 672 | ) |
| 673 | for browser in browsers |
| 674 | ), |
| 675 | ) |
| 676 | if signature != last_signature: |
| 677 | await self._emit_viewer_state( |
| 678 | sid, |
| 679 | context_id, |
| 680 | active_id, |
| 681 | browsers=browsers, |
| 682 | viewer_id=viewer_id, |
| 683 | state=state, |
| 684 | viewer_transport=viewer_transport, |
| 685 | ) |
| 686 | last_signature = signature |
| 687 | await asyncio.sleep(SNAPSHOT_STATE_POLL_SECONDS) |
| 688 | except asyncio.CancelledError: |
| 689 | raise |
| 690 | except Exception: |
| 691 | await asyncio.sleep(FRAME_RETRY_DELAY_SECONDS) |
| 692 | |
| 693 | @staticmethod |
| 694 | def _active_browser_id( |
| 695 | listing: dict[str, Any], |
| 696 | requested_browser_id: int | str | None, |
| 697 | ) -> int | str | None: |
| 698 | browsers = listing.get("browsers") or [] |
| 699 | browser_ids = {str(browser.get("id")) for browser in browsers} |
| 700 | requested_id = str(requested_browser_id or "") if requested_browser_id else "" |
| 701 | active_id = ( |
| 702 | requested_browser_id |
| 703 | if requested_id and requested_id in browser_ids |
| 704 | else listing.get("last_interacted_browser_id") |
| 705 | ) |
| 706 | if active_id and str(active_id) not in browser_ids: |
| 707 | active_id = None |
| 708 | if not active_id and browsers: |
| 709 | active_id = browsers[0].get("id") |
| 710 | return active_id |
| 711 | |
| 712 | async def _effective_viewer( |
| 713 | self, |
| 714 | runtime: Any, |
| 715 | browser_id: int | str | None, |
| 716 | data: dict[str, Any], |
| 717 | ) -> tuple[str, dict[str, Any] | None]: |
| 718 | requested = self._viewer_transport(data) |
| 719 | if requested != VIEWER_TRANSPORT_INTERACTIVE or not runtime or not browser_id: |
| 720 | return requested, None |
| 721 | viewport = self._viewport_from_data(data) or {} |
| 722 | try: |
| 723 | viewer = await runtime.call( |
| 724 | "interactive_viewer", |
| 725 | browser_id, |
| 726 | width=int(viewport.get("width") or 0), |
| 727 | height=int(viewport.get("height") or 0), |
| 728 | ) |
| 729 | except Exception as exc: |
| 730 | viewer = {"available": False, "error": str(exc)} |
| 731 | if viewer.get("available"): |
| 732 | return VIEWER_TRANSPORT_INTERACTIVE, viewer |
| 733 | return VIEWER_TRANSPORT_SCREENCAST, viewer |
| 734 | |
| 735 | @staticmethod |
| 736 | def _result_browser_id(result: Any) -> int | str | None: |
| 737 | if not isinstance(result, dict): |
| 738 | return None |
| 739 | state = result.get("state") if isinstance(result.get("state"), dict) else result |
| 740 | return state.get("id") if isinstance(state, dict) else None |
| 741 | |
| 742 | @staticmethod |
| 743 | def _state_for_browser( |
| 744 | browsers: list[dict[str, Any]], |
| 745 | browser_id: int | str, |
| 746 | current_state: dict[str, Any] | None, |
| 747 | ) -> dict[str, Any] | None: |
| 748 | for browser in browsers: |
| 749 | if str(browser.get("id")) == str(browser_id): |
| 750 | return browser |
| 751 | return current_state |
| 752 | |
| 753 | @staticmethod |
| 754 | def _state_signature( |
| 755 | active_id: int | str | None, |
| 756 | browsers: list[dict[str, Any]], |
| 757 | ) -> tuple[str, tuple[tuple[str, str, str, str, bool], ...]]: |
| 758 | return ( |
| 759 | str(active_id or ""), |
| 760 | tuple( |
| 761 | ( |
| 762 | str(browser.get("context_id") or ""), |
| 763 | str(browser.get("id") or ""), |
| 764 | str(browser.get("currentUrl") or ""), |
| 765 | str(browser.get("title") or ""), |
| 766 | bool(browser.get("loading")), |
| 767 | ) |
| 768 | for browser in browsers |
| 769 | ), |
| 770 | ) |
| 771 | |
| 772 | @staticmethod |
| 773 | def _frame_payload( |
| 774 | frame: dict[str, Any], |
| 775 | *, |
| 776 | context_id: str, |
| 777 | viewer_id: str, |
| 778 | browser_id: int | str, |
| 779 | sequence: int, |
| 780 | binary_frames: bool, |
| 781 | ) -> dict[str, Any]: |
| 782 | image = str(frame.get("image") or "") |
| 783 | payload: dict[str, Any] = { |
| 784 | "context_id": context_id, |
| 785 | "viewer_id": viewer_id, |
| 786 | "browser_id": browser_id, |
| 787 | "seq": sequence, |
| 788 | "mime": frame.get("mime") or "image/jpeg", |
| 789 | "frame_source": VIEWER_TRANSPORT_SCREENCAST, |
| 790 | "viewer_transport": VIEWER_TRANSPORT_SCREENCAST, |
| 791 | } |
| 792 | dimensions = WsBrowser._frame_dimensions(frame.get("metadata")) |
| 793 | if dimensions: |
| 794 | payload.update(dimensions) |
| 795 | if binary_frames: |
| 796 | try: |
| 797 | payload["image"] = base64.b64decode(image, validate=False) |
| 798 | payload["encoding"] = "binary" |
| 799 | except Exception: |
| 800 | payload["image"] = image |
| 801 | payload["encoding"] = "base64" |
| 802 | else: |
| 803 | payload["image"] = image |
| 804 | payload["encoding"] = "base64" |
| 805 | return payload |
| 806 | |
| 807 | @staticmethod |
| 808 | def _frame_dimensions(metadata: Any) -> dict[str, int]: |
| 809 | if not isinstance(metadata, dict): |
| 810 | return {} |
| 811 | |
| 812 | def dimensions(width_key: str, height_key: str) -> tuple[int, int] | None: |
| 813 | try: |
| 814 | width = int(metadata.get(width_key) or 0) |
| 815 | height = int(metadata.get(height_key) or 0) |
| 816 | except (TypeError, ValueError): |
| 817 | return None |
| 818 | if width > 0 and height > 0: |
| 819 | return width, height |
| 820 | return None |
| 821 | |
| 822 | expected = dimensions("expectedWidth", "expectedHeight") |
| 823 | jpeg = dimensions("jpegWidth", "jpegHeight") |
| 824 | if expected and jpeg: |
| 825 | width_scale = jpeg[0] / expected[0] |
| 826 | height_scale = jpeg[1] / expected[1] |
| 827 | if abs(width_scale - height_scale) <= 0.01: |
| 828 | return {"width": expected[0], "height": expected[1]} |
| 829 | return {"width": jpeg[0], "height": jpeg[1]} |
| 830 | |
| 831 | for fallback in (jpeg, dimensions("deviceWidth", "deviceHeight"), expected): |
| 832 | if fallback: |
| 833 | return {"width": fallback[0], "height": fallback[1]} |
| 834 | return {} |
| 835 | |
| 836 | async def _emit_viewer_state( |
| 837 | self, |
| 838 | sid: str, |
| 839 | context_id: str, |
| 840 | browser_id: int | str | None, |
| 841 | *, |
| 842 | browsers: list[dict[str, Any]] | None = None, |
| 843 | viewer_id: str = "", |
| 844 | state: dict[str, Any] | None = None, |
| 845 | viewer_transport: str = VIEWER_TRANSPORT_SNAPSHOT, |
| 846 | ) -> None: |
| 847 | browsers, all_browsers, tab_scope = await self._tabs_for_scope(context_id, browsers or []) |
| 848 | await self._emit_to_connected_viewer( |
| 849 | sid, |
| 850 | "browser_viewer_state", |
| 851 | { |
| 852 | "context_id": context_id, |
| 853 | "active_browser_context_id": context_id, |
| 854 | "viewer_id": viewer_id, |
| 855 | "browser_id": browser_id, |
| 856 | "active_browser_id": browser_id, |
| 857 | "browsers": browsers or [], |
| 858 | "state": state, |
| 859 | "all_browsers": all_browsers, |
| 860 | "tab_scope": tab_scope, |
| 861 | "viewer_transport": viewer_transport, |
| 862 | }, |
| 863 | ) |
| 864 | |
| 865 | async def _emit_empty_frame( |
| 866 | self, |
| 867 | sid: str, |
| 868 | context_id: str, |
| 869 | *, |
| 870 | browsers: list[dict[str, Any]] | None = None, |
| 871 | viewer_id: str = "", |
| 872 | frame_source: str = "", |
| 873 | ) -> None: |
| 874 | browsers, all_browsers, tab_scope = await self._tabs_for_scope(context_id, browsers or []) |
| 875 | await self._emit_to_connected_viewer( |
| 876 | sid, |
| 877 | "browser_viewer_frame", |
| 878 | { |
| 879 | "context_id": context_id, |
| 880 | "viewer_id": viewer_id, |
| 881 | "browser_id": None, |
| 882 | "browsers": browsers or [], |
| 883 | "all_browsers": all_browsers, |
| 884 | "tab_scope": tab_scope, |
| 885 | "image": "", |
| 886 | "mime": "", |
| 887 | "state": None, |
| 888 | "frame_source": frame_source, |
| 889 | "viewer_transport": frame_source or VIEWER_TRANSPORT_SNAPSHOT, |
| 890 | }, |
| 891 | ) |
| 892 | |
| 893 | async def _emit_to_connected_viewer( |
| 894 | self, |
| 895 | sid: str, |
| 896 | event: str, |
| 897 | data: dict[str, Any], |
| 898 | ) -> None: |
| 899 | manager = getattr(self, "_manager", None) |
| 900 | if manager is not None: |
| 901 | with manager.lock: |
| 902 | connected = (getattr(self, "namespace", "/ws"), sid) in manager.connections |
| 903 | if not connected: |
| 904 | raise asyncio.CancelledError() |
| 905 | await self.emit_to(sid, event, data) |
| 906 | |
| 907 | @staticmethod |
| 908 | def _viewer_transport(data: dict[str, Any]) -> str: |
| 909 | raw = ( |
| 910 | data.get("viewer_transport") |
| 911 | or data.get("viewerTransport") |
| 912 | or data.get("surface_transport") |
| 913 | or data.get("surfaceTransport") |
| 914 | or data.get("transport") |
| 915 | or "" |
| 916 | ) |
| 917 | normalized = str(raw or "").strip().lower().replace("-", "_") |
| 918 | if normalized == "live": |
| 919 | normalized = VIEWER_TRANSPORT_SCREENCAST |
| 920 | if normalized in VIEWER_TRANSPORTS: |
| 921 | return normalized |
| 922 | return VIEWER_TRANSPORT_SNAPSHOT |
| 923 | |
| 924 | @staticmethod |
| 925 | def _viewport_from_data(data: dict[str, Any]) -> dict[str, int] | None: |
| 926 | try: |
| 927 | width = int(data.get("viewport_width") or data.get("width") or 0) |
| 928 | height = int(data.get("viewport_height") or data.get("height") or 0) |
| 929 | except (TypeError, ValueError): |
| 930 | return None |
| 931 | if width < 80 or height < 80: |
| 932 | return None |
| 933 | return { |
| 934 | "width": max(320, min(4096, width)), |
| 935 | "height": max(200, min(4096, height)), |
| 936 | } |
| 937 | |
| 938 | @staticmethod |
| 939 | def _capture_scale_from_data(data: dict[str, Any]) -> float: |
| 940 | try: |
| 941 | scale = float( |
| 942 | data.get("device_pixel_ratio") |
| 943 | or data.get("devicePixelRatio") |
| 944 | or data.get("pixel_ratio") |
| 945 | or data.get("pixelRatio") |
| 946 | or 1 |
| 947 | ) |
| 948 | except (TypeError, ValueError): |
| 949 | return 1.0 |
| 950 | return max(1.0, min(2.0, scale)) |
| 951 | |
| 952 | @staticmethod |
| 953 | def _context_id(data: dict[str, Any]) -> str: |
| 954 | return str(data.get("context_id") or data.get("context") or "").strip() |
| 955 | |
| 956 | @staticmethod |
| 957 | def _bool(value: Any) -> bool: |
| 958 | if isinstance(value, bool): |
| 959 | return value |
| 960 | if isinstance(value, (int, float)): |
| 961 | return bool(value) |
| 962 | return str(value or "").strip().lower() in {"1", "true", "yes", "on"} |
| 963 | |
| 964 | @staticmethod |
| 965 | def _error(code: str, message: str, data: dict[str, Any]) -> WsResult: |
| 966 | return WsResult.error( |
| 967 | code=code, |
| 968 | message=message, |
| 969 | correlation_id=data.get("correlationId"), |
| 970 | ) |