main
py 970 lines 36.8 KB
Raw
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 )