From 4364f1328ffcee1080f7643583879fe3464e3110 Mon Sep 17 00:00:00 2001
From: suchintan <3853670+suchintan@users.noreply.github.com>
Date: Sat, 15 Aug 2026 18:45:10 +0000
Subject: [PATCH 1/7] =?UTF-8?q?=F0=9F=94=84=20synced=20local=20'benchmarks?=
=?UTF-8?q?/'=20with=20remote=20'benchmarks/'?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
https://github.com/Skyvern-AI/rustwright-cloud/pull/209
---
benchmarks/automation_cases.py | 8 ++++++--
1 file changed, 6 insertions(+), 2 deletions(-)
diff --git a/benchmarks/automation_cases.py b/benchmarks/automation_cases.py
index eb961b3..2d0b5f1 100644
--- a/benchmarks/automation_cases.py
+++ b/benchmarks/automation_cases.py
@@ -4292,7 +4292,9 @@ def page_event_waiters_reject_on_page_crash(page):
@case
def page_errors_history_and_clear(page):
page.set_content("page error history")
- page.evaluate("() => setTimeout(() => { throw new Error('parity page boom'); }, 0)")
+ page.add_script_tag(
+ content="setTimeout(() => { throw new Error('parity page boom'); }, 0)"
+ )
deadline = time.monotonic() + 3
errors = []
@@ -4313,7 +4315,9 @@ def page_errors_since_navigation_filter(page):
after = "page error after navigation filter"
after_set_content = "page error after set content filter"
page.set_content("page error filter")
- page.evaluate("(text) => setTimeout(() => { throw new Error(text); }, 0)", before)
+ page.add_script_tag(
+ content=f"setTimeout(() => {{ throw new Error({json.dumps(before)}); }}, 0)"
+ )
deadline = time.monotonic() + 3
while time.monotonic() < deadline:
From 1364f4a8599d15cfe7463acc19e14131586d2b21 Mon Sep 17 00:00:00 2001
From: suchintan <3853670+suchintan@users.noreply.github.com>
Date: Sat, 15 Aug 2026 18:45:10 +0000
Subject: [PATCH 2/7] =?UTF-8?q?=F0=9F=94=84=20synced=20local=20'python/'?=
=?UTF-8?q?=20with=20remote=20'python/'?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
https://github.com/Skyvern-AI/rustwright-cloud/pull/209
---
python/rustwright/_async_generated.py | 2 +-
python/rustwright/_compat/__init__.py | 204 +++++--
.../_compat/pytest_playwright/__init__.py | 11 +-
python/rustwright/sync_api.py | 508 ++++++++++++++++--
4 files changed, 620 insertions(+), 105 deletions(-)
diff --git a/python/rustwright/_async_generated.py b/python/rustwright/_async_generated.py
index 106c543..a6fcdf8 100644
--- a/python/rustwright/_async_generated.py
+++ b/python/rustwright/_async_generated.py
@@ -1,5 +1,5 @@
# This file is generated by tools/generate_async_api.py. Do not edit.
-# sync_api.py sha256: ee8829c2aee8bbcf37bb44950e754223ecc0f71492c5faedfb1dc48389827067
+# sync_api.py sha256: ea3c8280513eef43c9f177ffa60f88315945b86a549b42513aa9249933d28240
from __future__ import annotations
from pathlib import Path
diff --git a/python/rustwright/_compat/__init__.py b/python/rustwright/_compat/__init__.py
index c47eeb8..ae4c9ab 100644
--- a/python/rustwright/_compat/__init__.py
+++ b/python/rustwright/_compat/__init__.py
@@ -1,14 +1,29 @@
-"""Explicit opt-in Playwright/Patchright/Cloakbrowser import compatibility."""
+"""Explicit opt-in Playwright/Patchright/Cloakbrowser import compatibility.
+
+Aliases are installed eagerly. If pytest is not importable when
+:func:`enable_playwright_compat` runs, pytest-plugin aliases are skipped. Call
+the function again after pytest becomes importable to add them. Pytest users
+normally enable compatibility inside a process where pytest is importable.
+
+Target imports happen before alias publication. If enable fails, canonical
+``rustwright._compat.*`` modules and pytest imported during that phase stay in
+``sys.modules``. Complete rollback covers only legacy alias entries and their
+parent attributes; removing canonical imports could disturb unrelated users.
+
+Do not enable or disable compatibility concurrently with in-flight imports of
+aliased names. Direct ``sys.modules`` aliasing cannot make those imports atomic.
+"""
from __future__ import annotations
import importlib
+import importlib.util
import sys
+from threading import RLock
from types import ModuleType
-from typing import Optional
-
+from typing import NamedTuple
-_ALIASES = (
+_CORE_ALIASES = (
("playwright", "rustwright._compat.playwright"),
("playwright.__main__", "rustwright._compat.playwright.__main__"),
("playwright._impl", "rustwright._compat.playwright._impl"),
@@ -16,7 +31,6 @@
("playwright._impl._errors", "rustwright._compat.playwright._impl._errors"),
("playwright.async_api", "rustwright._compat.playwright.async_api"),
("playwright.async_api._generated", "rustwright._compat.playwright.async_api._generated"),
- ("playwright.pytest_plugin", "rustwright._compat.playwright.pytest_plugin"),
("playwright.sync_api", "rustwright._compat.playwright.sync_api"),
("playwright.sync_api._generated", "rustwright._compat.playwright.sync_api._generated"),
("patchright", "rustwright._compat.patchright"),
@@ -26,69 +40,163 @@
("patchright._impl._errors", "rustwright._compat.patchright._impl._errors"),
("patchright.async_api", "rustwright._compat.patchright.async_api"),
("patchright.async_api._generated", "rustwright._compat.patchright.async_api._generated"),
- ("patchright.pytest_plugin", "rustwright._compat.patchright.pytest_plugin"),
("patchright.sync_api", "rustwright._compat.patchright.sync_api"),
("patchright.sync_api._generated", "rustwright._compat.patchright.sync_api._generated"),
("cloakbrowser", "rustwright._compat.cloakbrowser"),
- # The pytest_playwright aliases re-export the full rustwright plugin. A
- # real pytest-playwright distribution's entry point resolving here loads
- # the plugin a second time, which is safe by construction: option
- # registration skips already-taken flags and browser_name parametrization
- # is guarded to run at most once per test.
+)
+
+_PYTEST_ALIASES = (
("pytest_playwright", "rustwright._compat.pytest_playwright"),
+ ("playwright.pytest_plugin", "rustwright._compat.playwright.pytest_plugin"),
+ ("patchright.pytest_plugin", "rustwright._compat.patchright.pytest_plugin"),
("pytest_playwright.pytest_playwright", "rustwright._compat.pytest_playwright.pytest_playwright"),
)
-_PREVIOUS_MODULES: dict[str, Optional[ModuleType]] = {}
+_MISSING = object()
+_PREVIOUS_MODULES: dict[str, object] = {}
+_PREVIOUS_PARENT_ATTRIBUTES: dict[str, tuple[ModuleType, object]] = {}
+_STATE_LOCK = RLock()
_ENABLED = False
+_PYTEST_ALIASES_ENABLED = False
-def _set_parent_attribute(module_name: str, module: ModuleType) -> None:
- parent_name, _, child_name = module_name.rpartition(".")
- if not parent_name:
- return
- parent = sys.modules.get(parent_name)
- if parent is not None:
- setattr(parent, child_name, module)
+class PlaywrightCompatEnableResult(NamedTuple):
+ """Aliases registered or skipped by the active compatibility state."""
+ enabled: bool
+ registered_aliases: tuple[str, ...]
+ skipped_aliases: tuple[str, ...]
-def enable_playwright_compat() -> None:
- """Enable legacy Playwright-compatible import names for this Python process.
- After this is called, subsequent imports such as ``playwright.sync_api`` or
- ``patchright.async_api`` resolve to Rustwright's compatibility shims.
- """
+_LAST_ENABLE_RESULT = PlaywrightCompatEnableResult(False, (), ())
- global _ENABLED
- if _ENABLED:
- return
- loaded_modules = [(alias_name, importlib.import_module(target_name)) for alias_name, target_name in _ALIASES]
- for alias_name, module in loaded_modules:
- _PREVIOUS_MODULES[alias_name] = sys.modules.get(alias_name)
- sys.modules[alias_name] = module
- _set_parent_attribute(alias_name, module)
+def _compat_transaction_hook(event: str, alias_name: str | None = None) -> None:
+ """Stable no-op event seam for compatibility transaction tests."""
- _ENABLED = True
-
-def disable_playwright_compat() -> None:
- """Undo aliases installed by :func:`enable_playwright_compat`."""
-
- global _ENABLED
- if not _ENABLED:
+def _set_parent_attribute(
+ module_name: str,
+ module: ModuleType,
+ previous_parent_attributes: dict[str, tuple[ModuleType, object]],
+) -> None:
+ parent_name, _, child_name = module_name.rpartition(".")
+ parent = sys.modules.get(parent_name)
+ if not isinstance(parent, ModuleType):
return
-
- for alias_name, _target_name in _ALIASES:
- previous = _PREVIOUS_MODULES.get(alias_name)
- if previous is None:
+ if module_name not in previous_parent_attributes:
+ previous_parent_attributes[module_name] = (
+ parent,
+ vars(parent).get(child_name, _MISSING),
+ )
+ ModuleType.__setattr__(parent, child_name, module)
+
+
+def _restore_aliases(
+ previous_modules: dict[str, object],
+ previous_parent_attributes: dict[str, tuple[ModuleType, object]],
+) -> None:
+ for alias_name in sorted(previous_modules, key=lambda name: name.count("."), reverse=True):
+ previous_module = previous_modules[alias_name]
+ if previous_module is _MISSING:
sys.modules.pop(alias_name, None)
else:
- sys.modules[alias_name] = previous
- _set_parent_attribute(alias_name, previous)
+ sys.modules[alias_name] = previous_module
+
+ parent_snapshot = previous_parent_attributes.get(alias_name)
+ if parent_snapshot is None:
+ continue
+ parent, previous_attribute = parent_snapshot
+ child_name = alias_name.rpartition(".")[2]
+ if previous_attribute is _MISSING:
+ if child_name in vars(parent):
+ ModuleType.__delattr__(parent, child_name)
+ else:
+ ModuleType.__setattr__(parent, child_name, previous_attribute)
- _PREVIOUS_MODULES.clear()
- _ENABLED = False
+def enable_playwright_compat() -> PlaywrightCompatEnableResult:
+ """Enable legacy aliases and report any aliases skipped without pytest.
-__all__ = ["disable_playwright_compat", "enable_playwright_compat"]
+ Every target import completes before the compatibility lock is acquired.
+ Calling again after pytest becomes importable upgrades an enabled core-only
+ state with the pytest aliases.
+
+ Target imports are outside the rollback boundary. Successfully imported
+ canonical compatibility modules, including pytest dependencies, remain in
+ ``sys.modules`` if enable later fails. Legacy alias entries and their parent
+ attributes are the complete transactional publication boundary.
+ """
+
+ global _ENABLED, _LAST_ENABLE_RESULT, _PYTEST_ALIASES_ENABLED
+
+ pytest_available = importlib.util.find_spec("pytest") is not None
+ aliases_to_import = _CORE_ALIASES + (_PYTEST_ALIASES if pytest_available else ())
+ _compat_transaction_hook("enable-before-import")
+ loaded_modules = tuple(
+ (alias_name, importlib.import_module(target_name))
+ for alias_name, target_name in aliases_to_import
+ )
+ _compat_transaction_hook("enable-after-import")
+
+ with _STATE_LOCK:
+ _compat_transaction_hook("enable-lock-acquired")
+ if _ENABLED:
+ if not pytest_available or _PYTEST_ALIASES_ENABLED:
+ return _LAST_ENABLE_RESULT
+ modules_to_publish = loaded_modules[len(_CORE_ALIASES) :]
+ else:
+ modules_to_publish = loaded_modules
+
+ previous_modules = {
+ alias_name: sys.modules.get(alias_name, _MISSING)
+ for alias_name, _module in modules_to_publish
+ }
+ previous_parent_attributes: dict[str, tuple[ModuleType, object]] = {}
+ try:
+ for alias_name, module in modules_to_publish:
+ sys.modules[alias_name] = module
+ _set_parent_attribute(alias_name, module, previous_parent_attributes)
+ _compat_transaction_hook("enable-after-alias-publish", alias_name)
+ except BaseException:
+ _restore_aliases(previous_modules, previous_parent_attributes)
+ raise
+
+ _PREVIOUS_MODULES.update(previous_modules)
+ _PREVIOUS_PARENT_ATTRIBUTES.update(previous_parent_attributes)
+ _ENABLED = True
+ _PYTEST_ALIASES_ENABLED = _PYTEST_ALIASES_ENABLED or pytest_available
+ skipped_aliases = (
+ ()
+ if _PYTEST_ALIASES_ENABLED
+ else tuple(alias_name for alias_name, _target_name in _PYTEST_ALIASES)
+ )
+ _LAST_ENABLE_RESULT = PlaywrightCompatEnableResult(
+ True,
+ tuple(_PREVIOUS_MODULES),
+ skipped_aliases,
+ )
+ return _LAST_ENABLE_RESULT
+
+
+def disable_playwright_compat() -> None:
+ """Restore modules and parent attributes replaced by compatibility."""
+
+ global _ENABLED, _LAST_ENABLE_RESULT, _PYTEST_ALIASES_ENABLED
+
+ with _STATE_LOCK:
+ if not _ENABLED:
+ return
+ _restore_aliases(_PREVIOUS_MODULES, _PREVIOUS_PARENT_ATTRIBUTES)
+ _PREVIOUS_MODULES.clear()
+ _PREVIOUS_PARENT_ATTRIBUTES.clear()
+ _ENABLED = False
+ _PYTEST_ALIASES_ENABLED = False
+ _LAST_ENABLE_RESULT = PlaywrightCompatEnableResult(False, (), ())
+
+
+__all__ = [
+ "PlaywrightCompatEnableResult",
+ "disable_playwright_compat",
+ "enable_playwright_compat",
+]
diff --git a/python/rustwright/_compat/pytest_playwright/__init__.py b/python/rustwright/_compat/pytest_playwright/__init__.py
index 82b3f29..8e39df2 100644
--- a/python/rustwright/_compat/pytest_playwright/__init__.py
+++ b/python/rustwright/_compat/pytest_playwright/__init__.py
@@ -1,3 +1,10 @@
-from .pytest_playwright import CreateContextCallback
+"""Compatibility surface for the optional pytest plugin."""
-__all__ = ["CreateContextCallback"]
+from __future__ import annotations
+
+from typing import TYPE_CHECKING
+
+from .pytest_playwright import *
+
+if TYPE_CHECKING:
+ from .pytest_playwright import CreateContextCallback as CreateContextCallback
diff --git a/python/rustwright/sync_api.py b/python/rustwright/sync_api.py
index 835b8ce..5c2e2db 100644
--- a/python/rustwright/sync_api.py
+++ b/python/rustwright/sync_api.py
@@ -3092,12 +3092,12 @@ def _file_chooser_upload_paths(files: Any, temporary_directories: list[tempfile.
def _contains_file_payload(files: Any) -> bool:
- if isinstance(files, dict):
+ if isinstance(files, (dict, bytes, bytearray)):
return True
- if files is None or isinstance(files, (str, Path, bytes, bytearray)):
+ if files is None or isinstance(files, (str, Path)):
return False
try:
- return any(isinstance(item, dict) for item in files)
+ return any(isinstance(item, (dict, bytes, bytearray)) for item in files)
except TypeError:
return False
@@ -3407,6 +3407,48 @@ def _wait_for_url_timeout_error(timeout_ms: float) -> TimeoutError:
def _method_timeout_error(method: str, timeout_ms: float) -> TimeoutError:
return TimeoutError(f"{method}: Timeout {_format_timeout_value(timeout_ms)}ms exceeded.")
+def _file_chooser_deadline(timeout_ms: float) -> float:
+ return float("inf") if timeout_ms == 0 else time.monotonic() + timeout_ms / 1000
+
+
+def _file_chooser_remaining_ms(deadline: float, timeout_ms: float) -> float:
+ if math.isinf(deadline):
+ return 0.0
+ remaining_ms = (deadline - time.monotonic()) * 1000
+ if remaining_ms <= 0:
+ raise _method_timeout_error("FileChooser.set_files", timeout_ms)
+ return max(remaining_ms, 1.0)
+
+
+def _file_chooser_cleanup_timeout_ms(deadline: float) -> Optional[float]:
+ if math.isinf(deadline):
+ return 1_000.0
+ remaining_ms = (deadline - time.monotonic()) * 1000
+ # Rust floors positive sub-millisecond command timeouts to 1 ms. Skip the
+ # release instead of exceeding the operation deadline; teardown reclaims it.
+ if remaining_ms < 1.0:
+ return None
+ return min(remaining_ms, 1_000.0)
+
+
+def _file_chooser_directory_inventory(
+ directories: list[Path],
+ *,
+ deadline: float,
+ timeout_ms: float,
+) -> list[str]:
+ expected_files: list[str] = []
+ for directory in directories:
+ _file_chooser_remaining_ms(deadline, timeout_ms)
+ for child in directory.rglob("*"):
+ _file_chooser_remaining_ms(deadline, timeout_ms)
+ is_file = child.is_file()
+ _file_chooser_remaining_ms(deadline, timeout_ms)
+ if is_file:
+ expected_files.append(f"{directory.name}/{child.relative_to(directory).as_posix()}")
+ _file_chooser_remaining_ms(deadline, timeout_ms)
+ return expected_files
+
def _resolve_url_match_base(base_url: Optional[str], expected: Any) -> Any:
if not base_url or not isinstance(expected, str) or expected.startswith("*"):
@@ -9637,8 +9679,22 @@ def __init__(self, page: Optional["Page"], payload: dict[str, Any], *, worker: O
self.type = str(payload.get("type") or "log")
self.text = str(payload.get("text") or "")
self.timestamp = float(payload.get("timestamp") or time.time() * 1000)
+ owner_frame = None
+ if page is not None and worker is None:
+ session_id = payload.get("session_id")
+ execution_context_id = payload.get("execution_context_id")
+ if session_id is not None and execution_context_id is not None:
+ try:
+ frame_id = page._core.execution_context_frame_id(
+ str(session_id),
+ str(execution_context_id),
+ )
+ if frame_id:
+ owner_frame = page._frame_for_id(str(frame_id))
+ except (AttributeError, Error):
+ owner_frame = None
self.args = [
- JSHandle(page or worker, _console_arg_handle_payload(arg))
+ JSHandle(page or worker, _console_arg_handle_payload(arg), owner_frame=owner_frame)
for arg in list(payload.get("args") or [])
]
location = payload.get("location")
@@ -9732,6 +9788,7 @@ def __init__(self, page: "Page", payload: dict[str, Any]):
self._frame_id = str(payload.get("frame_id") or "")
self._backend_node_id = int(payload.get("backend_node_id") or 0)
self._mode = str(payload.get("mode") or "")
+ self._session_id = str(payload.get("session_id") or "")
self._temporary_upload_dirs: list[tempfile.TemporaryDirectory[str]] = []
def is_multiple(self) -> bool:
@@ -9741,59 +9798,228 @@ def is_multiple(self) -> bool:
def page(self) -> "Page":
return self._page
+ def _resolve_backend_node_element(
+ self,
+ *,
+ method: str,
+ command_timeout_ms: float,
+ mutation: bool,
+ ) -> "ElementHandle":
+ core_session = (
+ _call_with_method_prefix(
+ method,
+ self._page._core.cdp_session_for_id,
+ self._session_id,
+ )
+ if self._session_id
+ else _call_with_method_prefix(method, self._page._core.cdp_session)
+ )
+ payload = json.loads(
+ _call_with_method_prefix(
+ method,
+ core_session.send,
+ "DOM.resolveNode",
+ json.dumps({"backendNodeId": self._backend_node_id}),
+ command_timeout_ms,
+ )
+ )
+ remote = payload.get("object")
+ if not isinstance(remote, dict) or not remote.get("objectId"):
+ raise Error(f"{method}: Element is not attached to the DOM")
+ class_name = remote.get("className")
+ if remote.get("subtype") != "node" and not (
+ isinstance(class_name, str) and class_name.endswith("Element")
+ ):
+ raise Error(f"{method}: JSHandle is not an Element")
+ if self._session_id:
+ remote["__rustwright_session_id"] = self._session_id
+ if self._frame_id:
+ remote["__rustwright_realm_identity"] = f"frame:{self._frame_id}"
+ owner_frame = (
+ self._page.main_frame
+ if mutation or not self._frame_id
+ else self._page._frame_for_id(self._frame_id)
+ )
+ handle = JSHandle(self._page, remote, owner_frame=owner_frame)
+ return ElementHandle(owner_frame.locator("*").nth(0), handle=handle)
+
+ def _mutation_element(self, *, deadline: float, timeout_ms: float) -> "ElementHandle":
+ command_timeout_ms = _file_chooser_remaining_ms(deadline, timeout_ms)
+ try:
+ if self._backend_node_id:
+ return self._resolve_backend_node_element(
+ method="FileChooser.set_files",
+ command_timeout_ms=command_timeout_ms,
+ mutation=True,
+ )
+ locator = self._page.locator("input[type=file]").first
+ handle = locator._evaluate_handle_with_method(
+ "(element) => element",
+ timeout=command_timeout_ms,
+ method="FileChooser.set_files",
+ )
+ return ElementHandle(locator, handle=handle)
+ except TimeoutError:
+ raise _method_timeout_error("FileChooser.set_files", timeout_ms) from None
+
+ def _dispose_mutation_element(self, element: "ElementHandle", *, deadline: float) -> None:
+ cleanup_timeout_ms = _file_chooser_cleanup_timeout_ms(deadline)
+ if cleanup_timeout_ms is None:
+ return
+ try:
+ element._dispose_with_timeout(cleanup_timeout_ms)
+ except Exception:
+ # This release is best-effort and must not replace the set_files result.
+ pass
+
@property
def element(self) -> "ElementHandle":
if self._backend_node_id:
- handle: Optional[JSHandle] = None
try:
- session = CDPSession(_call(self._page._core.cdp_session))
- payload = session.send("DOM.resolveNode", {"backendNodeId": self._backend_node_id})
- remote = payload.get("object")
- if isinstance(remote, dict):
- handle = JSHandle(self._page, remote)
- element = handle.as_element()
- if element is not None:
- handle = None
- return element
+ return self._resolve_backend_node_element(
+ method="FileChooser.element",
+ command_timeout_ms=30_000.0,
+ mutation=False,
+ )
except Exception:
pass
- finally:
- if handle is not None:
- handle.dispose()
- handle = _element_handle_from_locator(self._page.locator("input[type=file]").first)
- return handle
+ return _element_handle_from_locator(self._page.locator("input[type=file]").first)
def set_files(self, files: Any, *, timeout: Optional[float] = None, no_wait_after: Optional[bool] = None) -> None:
timeout_ms = _default_timeout_for_method(self._page, timeout, method="FileChooser.set_files")
+ deadline = _file_chooser_deadline(timeout_ms)
upload_files = files
if files is not None and not isinstance(files, (str, Path, bytes, bytearray, dict)):
try:
upload_files = list(files)
except TypeError:
upload_files = files
+ _file_chooser_remaining_ms(deadline, timeout_ms)
if _contains_file_payload(upload_files):
payloads = _file_payloads(upload_files)
+ _file_chooser_remaining_ms(deadline, timeout_ms)
if len(payloads) > 1 and not self.is_multiple():
raise Error("FileChooser.set_files: Error: Non-multiple file input can only accept single file")
- element = self.element
- locator = element._live_locator("set_input_files")
- assert locator is not None
- locator._set_input_files_impl(
- "FileChooser.set_files",
- upload_files,
- timeout=timeout_ms,
- no_wait_after=no_wait_after,
- )
+ element = self._mutation_element(deadline=deadline, timeout_ms=timeout_ms)
+ try:
+ try:
+ element._evaluate_with_timeout(
+ """(input, payloads) => {
+ if (!(input instanceof HTMLInputElement) || input.type !== 'file')
+ throw new Error('Element is not a file input');
+ const transfer = new DataTransfer();
+ for (const payload of payloads) {
+ const binary = atob(payload.buffer || '');
+ const bytes = new Uint8Array(binary.length);
+ for (let index = 0; index < binary.length; index++)
+ bytes[index] = binary.charCodeAt(index);
+ transfer.items.add(new File(
+ [bytes],
+ payload.name || 'file',
+ { type: payload.mime_type || '' },
+ ));
+ }
+ input.files = transfer.files;
+ input.dispatchEvent(new Event('input', { bubbles: true }));
+ input.dispatchEvent(new Event('change', { bubbles: true }));
+ }""",
+ payloads,
+ timeout_ms=_file_chooser_remaining_ms(deadline, timeout_ms),
+ method="FileChooser.set_files",
+ )
+ except TimeoutError:
+ raise _method_timeout_error("FileChooser.set_files", timeout_ms) from None
+ finally:
+ self._dispose_mutation_element(element, deadline=deadline)
return
paths = _file_chooser_upload_paths(upload_files, self._temporary_upload_dirs)
+ _file_chooser_remaining_ms(deadline, timeout_ms)
if len(paths) > 1 and not self.is_multiple():
raise Error("FileChooser.set_files: Error: Non-multiple file input can only accept single file")
- _call(
- self._page._core.set_file_input_files,
- self._backend_node_id,
- json_module_dumps(paths),
- timeout_ms,
- )
+ if not paths:
+ element = self._mutation_element(deadline=deadline, timeout_ms=timeout_ms)
+ try:
+ try:
+ element._evaluate_with_timeout(
+ """input => {
+ input.files = new DataTransfer().files;
+ input.dispatchEvent(new Event('input', { bubbles: true }));
+ input.dispatchEvent(new Event('change', { bubbles: true }));
+ }""",
+ timeout_ms=_file_chooser_remaining_ms(deadline, timeout_ms),
+ method="FileChooser.set_files",
+ )
+ except TimeoutError:
+ raise _method_timeout_error("FileChooser.set_files", timeout_ms) from None
+ finally:
+ self._dispose_mutation_element(element, deadline=deadline)
+ return
+ try:
+ _call_with_method_prefix(
+ "FileChooser.set_files",
+ self._page._core.set_file_input_files,
+ self._backend_node_id,
+ json_module_dumps(paths),
+ _file_chooser_remaining_ms(deadline, timeout_ms),
+ self._session_id or None,
+ )
+ except TimeoutError:
+ raise _method_timeout_error("FileChooser.set_files", timeout_ms) from None
+ directories: list[Path] = []
+ for path in paths:
+ _file_chooser_remaining_ms(deadline, timeout_ms)
+ candidate = Path(path)
+ is_directory = candidate.is_dir()
+ _file_chooser_remaining_ms(deadline, timeout_ms)
+ if is_directory:
+ directories.append(candidate)
+ if directories:
+ expected_files = _file_chooser_directory_inventory(
+ directories,
+ deadline=deadline,
+ timeout_ms=timeout_ms,
+ )
+ ready_expression = """(input, expected) => {
+ const actual = new Map();
+ for (const file of input.files) {
+ const path = file.webkitRelativePath;
+ actual.set(path, (actual.get(path) || 0) + 1);
+ }
+ return input.files.length === expected.length && expected.every(value => {
+ const count = actual.get(value) || 0;
+ if (!count) return false;
+ actual.set(value, count - 1);
+ return true;
+ });
+ }"""
+ else:
+ expected_files = []
+ for path in paths:
+ _file_chooser_remaining_ms(deadline, timeout_ms)
+ expected_files.append(Path(path).name)
+ ready_expression = """(input, expected) => {
+ const actual = Array.from(input.files, file => file.name);
+ return actual.length === expected.length
+ && actual.every((value, index) => value === expected[index]);
+ }"""
+ element = self._mutation_element(deadline=deadline, timeout_ms=timeout_ms)
+ try:
+ while True:
+ try:
+ ready = element._evaluate_with_timeout(
+ ready_expression,
+ expected_files,
+ timeout_ms=_file_chooser_remaining_ms(deadline, timeout_ms),
+ method="FileChooser.set_files",
+ )
+ except TimeoutError:
+ raise _method_timeout_error("FileChooser.set_files", timeout_ms) from None
+ if ready:
+ return
+ _file_chooser_remaining_ms(deadline, timeout_ms)
+ _sleep_until_next_poll(deadline)
+ finally:
+ self._dispose_mutation_element(element, deadline=deadline)
class JSHandle(_EventEmitter):
@@ -9802,12 +10028,31 @@ def __init__(self, page: Any, payload: dict[str, Any], *, owner_frame: Optional[
self._payload = payload
self._object_id = payload.get("objectId")
self._session_id = payload.get("__rustwright_session_id")
+ self._realm_identity_override = payload.get("__rustwright_realm_identity")
self._owner_frame = owner_frame
self._disposed = False
def _session_args(self) -> tuple[str, ...]:
return () if self._session_id is None else (str(self._session_id),)
+ def _realm_identity(self) -> Optional[str]:
+ if self._realm_identity_override is not None:
+ return str(self._realm_identity_override)
+ if self._owner_frame is not None:
+ frame_id = getattr(self._owner_frame, "_frame_id", None)
+ if frame_id:
+ return f"frame:{frame_id}"
+ target_id = getattr(self._page, "_target_id", None)
+ if target_id:
+ return f"worker:{target_id}"
+ return None
+
+ def _serialized_owner_args(self) -> tuple[Any, ...]:
+ realm_identity = self._realm_identity()
+ if self._session_id is None:
+ return () if self._owner_frame is None else (None, realm_identity)
+ return (str(self._session_id), realm_identity)
+
def _preview(self) -> str:
if self._payload.get("subtype") == "node":
return "JSHandle@node"
@@ -9848,7 +10093,7 @@ def json_value(self) -> Any:
self._page._core.js_handle_json_value,
self._object_id,
self._page._default_timeout,
- *self._session_args(),
+ *self._serialized_owner_args(),
)
)
)
@@ -9871,7 +10116,7 @@ def _truthy(self, timeout_ms: Optional[float] = None) -> bool:
None,
True,
self._page._default_timeout if timeout_ms is None else timeout_ms,
- *self._session_args(),
+ *self._serialized_owner_args(),
)
return bool(_decode_json_result(json.loads(result)))
if self._payload.get("type") == "undefined" or self._payload.get("subtype") == "null":
@@ -9950,8 +10195,17 @@ def as_element(self) -> Optional["ElementHandle"]:
except Error:
return None
- def _evaluate_with_method(self, expression: str, arg: Any = None, *, method: str) -> Any:
+ def _evaluate_with_method(
+ self,
+ expression: str,
+ arg: Any = None,
+ *,
+ method: str,
+ timeout_ms: Optional[float] = None,
+ ) -> Any:
expression = _normalize_string_option(expression, method=method, name="expression")
+ if arg is not None:
+ _ensure_evaluate_argument_context(self._page, self._owner_frame, arg, method=method)
if hasattr(self._page, "_mark_history_events_may_arrive"):
self._page._mark_history_events_may_arrive()
if not self._object_id:
@@ -9976,8 +10230,8 @@ def _evaluate_with_method(self, expression: str, arg: Any = None, *, method: str
_evaluate_handle_argument_function(expression),
json_module_dumps(prepared.cdp_arguments()),
True,
- None,
- *self._session_args(),
+ timeout_ms,
+ *self._serialized_owner_args(),
)
return _decode_json_result(json.loads(result))
finally:
@@ -9991,8 +10245,8 @@ def _evaluate_with_method(self, expression: str, arg: Any = None, *, method: str
expression,
arg_json,
True,
- None,
- *self._session_args(),
+ timeout_ms,
+ *self._serialized_owner_args(),
)
return _decode_json_result(json.loads(result))
@@ -10002,6 +10256,8 @@ def evaluate(self, expression: str, arg: Any = None) -> Any:
def _evaluate_handle_with_method(self, expression: str, arg: Any = None, *, method: str) -> "JSHandle":
expression = _normalize_string_option(expression, method=method, name="expression")
+ if arg is not None:
+ _ensure_evaluate_argument_context(self._page, self._owner_frame, arg, method=method)
if hasattr(self._page, "_mark_history_events_may_arrive"):
self._page._mark_history_events_may_arrive()
if not self._object_id:
@@ -10028,7 +10284,7 @@ def _evaluate_handle_with_method(self, expression: str, arg: Any = None, *, meth
json_module_dumps(prepared.cdp_arguments()),
False,
None,
- *self._session_args(),
+ *self._serialized_owner_args(),
)
)
return JSHandle(self._page, payload, owner_frame=self._owner_frame)
@@ -10045,7 +10301,7 @@ def _evaluate_handle_with_method(self, expression: str, arg: Any = None, *, meth
arg_json,
False,
None,
- *self._session_args(),
+ *self._serialized_owner_args(),
)
)
return JSHandle(self._page, payload, owner_frame=self._owner_frame)
@@ -10054,18 +10310,21 @@ def evaluate_handle(self, expression: str, arg: Any = None) -> "JSHandle":
self._ensure_not_disposed("evaluate_handle")
return self._evaluate_handle_with_method(expression, arg, method="JSHandle.evaluate_handle")
- def dispose(self) -> None:
+ def _dispose_with_timeout(self, timeout_ms: Optional[float]) -> None:
if self._disposed:
return
if self._object_id:
_call(
self._page._core.js_handle_dispose,
self._object_id,
- None,
+ timeout_ms,
*self._session_args(),
)
self._disposed = True
+ def dispose(self) -> None:
+ self._dispose_with_timeout(None)
+
def json_module_dumps(value: Any) -> str:
return __import__("json").dumps(value, separators=(",", ":"))
@@ -13478,6 +13737,10 @@ def __init__(
self._uses_direct_evaluation = False
self._child_frame_cache: list["Frame"] = []
+ def _raise_if_detached(self, method: str) -> None:
+ if self.is_detached():
+ raise Error(f"{method}: Frame was detached")
+
def _remember_child_frame(self, frame: "Frame") -> None:
if all(existing is not frame for existing in self._child_frame_cache):
self._child_frame_cache.append(frame)
@@ -13774,7 +14037,15 @@ def wait_for_selector(
return _element_handle_from_locator(locator.nth(0)) if attached else None
def evaluate(self, expression: str, arg: Any = None) -> Any:
+ self._raise_if_detached("Frame.evaluate")
expression = _normalize_string_option(expression, method="Frame.evaluate", name="expression")
+ if arg is not None:
+ _ensure_evaluate_argument_context(
+ self._page,
+ self,
+ arg,
+ method="Frame.evaluate",
+ )
self._page._mark_request_cookie_sync_required()
self._page._mark_history_events_may_arrive()
if self._is_main:
@@ -13808,6 +14079,23 @@ def evaluate(self, expression: str, arg: Any = None) -> Any:
if result is not None:
self._uses_direct_evaluation = True
return _decode_json_result(json.loads(result))
+ if arg is not None and _argument_contains_handle(arg):
+ prepared = _prepare_evaluate_argument(self._page, arg)
+ try:
+ anchor = prepared.handles[0]
+ result = _call_with_method_prefix(
+ "Frame.evaluate",
+ self._page._core.js_handle_evaluate_with_call_arguments,
+ anchor._object_id,
+ _evaluate_argument_function(expression),
+ json_module_dumps(prepared.cdp_arguments()),
+ True,
+ None,
+ *anchor._serialized_owner_args(),
+ )
+ return _decode_json_result(json.loads(result))
+ finally:
+ prepared.dispose_temporaries()
if self._frame_spec is not None:
return self.locator(":root")._evaluate_with_method(
"""(el, payload) => {
@@ -13860,11 +14148,38 @@ def evaluate(self, expression: str, arg: Any = None) -> Any:
)
def evaluate_handle(self, expression: str, arg: Any = None) -> JSHandle:
+ self._raise_if_detached("Frame.evaluate_handle")
expression = _normalize_string_option(expression, method="Frame.evaluate_handle", name="expression")
+ if arg is not None:
+ _ensure_evaluate_argument_context(
+ self._page,
+ self,
+ arg,
+ method="Frame.evaluate_handle",
+ )
self._page._mark_request_cookie_sync_required()
self._page._mark_history_events_may_arrive()
if self._is_main:
return self._page._evaluate_handle_with_timeout(expression, arg, method="Frame.evaluate_handle")
+ if arg is not None and _argument_contains_handle(arg):
+ prepared = _prepare_evaluate_argument(self._page, arg)
+ try:
+ anchor = prepared.handles[0]
+ payload = json.loads(
+ _call_with_method_prefix(
+ "Frame.evaluate_handle",
+ self._page._core.js_handle_evaluate_with_call_arguments,
+ anchor._object_id,
+ _evaluate_argument_function(expression),
+ json_module_dumps(prepared.cdp_arguments()),
+ False,
+ None,
+ *anchor._serialized_owner_args(),
+ )
+ )
+ return JSHandle(self._page, payload, owner_frame=self)
+ finally:
+ prepared.dispose_temporaries()
handle = self.locator(":root")._evaluate_handle_with_method(
"""(el, payload) => {
const [expression, arg] = payload;
@@ -17503,6 +17818,8 @@ def emulate_media(
def _evaluate_with_method(self, expression: str, arg: Any = None, *, method: str) -> Any:
expression = _normalize_string_option(expression, method=method, name="expression")
+ if arg is not None:
+ _ensure_evaluate_argument_context(self, self._main_frame, arg, method=method)
self._mark_request_cookie_sync_required()
self._mark_history_events_may_arrive()
console_marker = self._console_dispatch_marker_if_listening()
@@ -17562,6 +17879,8 @@ def _evaluate_handle_with_timeout(
method: str = "Page.evaluate_handle",
) -> JSHandle:
expression = _normalize_string_option(expression, method=method, name="expression")
+ if arg is not None:
+ _ensure_evaluate_argument_context(self, self._main_frame, arg, method=method)
self._mark_request_cookie_sync_required()
self._mark_history_events_may_arrive()
command_timeout = timeout_ms
@@ -17578,7 +17897,7 @@ def _evaluate_handle_with_timeout(
command_timeout,
)
)
- return JSHandle(self, payload)
+ return JSHandle(self, payload, owner_frame=self._main_frame)
finally:
prepared.dispose_temporaries()
arg = prepared.value if prepared is not None else arg
@@ -17592,7 +17911,7 @@ def _evaluate_handle_with_timeout(
command_timeout,
)
)
- return JSHandle(self, payload)
+ return JSHandle(self, payload, owner_frame=self._main_frame)
def evaluate_handle(self, expression: str, arg: Any = None) -> JSHandle:
return self._evaluate_handle_with_timeout(expression, arg)
@@ -21690,6 +22009,12 @@ def _wait_for_fill_ready(self, action: str, *, timeout: Optional[float] = None)
f"timed out waiting for locator to be editable while trying to {action}; {detail}"
)
+ def _owner_frame(self) -> Frame:
+ frame_spec = _frame_scope_spec_from_element_spec(self._spec)
+ if frame_spec is None:
+ return self._page.main_frame
+ return self._page._frame_from_spec(frame_spec)
+
def _evaluate_with_method(
self,
expression: str,
@@ -21757,7 +22082,7 @@ def _evaluate_handle_with_method(
self._page._default_timeout if timeout is None else timeout,
)
)
- return JSHandle(self._page, payload)
+ return JSHandle(self._page, payload, owner_frame=self._owner_frame())
def evaluate_handle(self, expression: str, arg: Any = None, *, timeout: Optional[float] = None) -> JSHandle:
return self._evaluate_handle_with_method(expression, arg, timeout=timeout, method="Locator.evaluate_handle")
@@ -24687,6 +25012,29 @@ def dispatch_event(self, type: str, event_init: Optional[dict[str, Any]] = None)
raise TargetClosedError(f"ElementHandle.dispatch_event: {_TARGET_CLOSED_MESSAGE}")
self._locator.dispatch_event(type, event_init)
+ def _evaluate_with_timeout(
+ self,
+ expression: str,
+ arg: Any = None,
+ *,
+ timeout_ms: float,
+ method: str,
+ ) -> Any:
+ self._ensure_not_disposed("evaluate")
+ if self._handle is not None and not self._handle._disposed:
+ return self._handle._evaluate_with_method(
+ expression,
+ arg,
+ method=method,
+ timeout_ms=timeout_ms,
+ )
+ return self._locator._evaluate_with_method(
+ expression,
+ arg,
+ timeout=timeout_ms,
+ method=method,
+ )
+
def evaluate(self, expression: str, arg: Any = None) -> Any:
self._ensure_not_disposed("evaluate")
if self._handle is not None and not self._handle._disposed:
@@ -25208,16 +25556,16 @@ def content_frame(self) -> Optional[Frame]:
return frame
def owner_frame(self) -> Frame:
- frame_spec = _frame_scope_spec_from_element_spec(self._locator._spec)
- if frame_spec is None:
- return self._locator._page.main_frame
- return self._locator._page._frame_from_spec(frame_spec)
+ return self._locator._owner_frame()
- def dispose(self) -> None:
+ def _dispose_with_timeout(self, timeout_ms: Optional[float]) -> None:
if self._handle is not None:
- self._handle.dispose()
+ self._handle._dispose_with_timeout(timeout_ms)
self._disposed = True
+ def dispose(self) -> None:
+ self._dispose_with_timeout(None)
+
_EVALUATE_HANDLE_MARKER = "__rustwright_handle_index__"
@@ -25243,6 +25591,54 @@ def dispose_temporaries(self) -> None:
pass
+def _same_evaluate_frame(left: Any, right: Any) -> bool:
+ if left is right:
+ return True
+ if left is None or right is None:
+ return False
+ left_id = getattr(left, "_frame_id", None)
+ right_id = getattr(right, "_frame_id", None)
+ return bool(left_id and right_id and left_id == right_id)
+
+
+def _ensure_evaluate_argument_context(
+ expected_owner: Any,
+ expected_frame: Optional[Any],
+ arg: Any,
+ *,
+ method: str,
+) -> None:
+ def validate(value: Any) -> None:
+ handle: Optional[JSHandle] = None
+ actual_owner: Any = None
+ owner_frame: Optional[Any] = None
+ if isinstance(value, JSHandle):
+ handle = value
+ actual_owner = value._page
+ owner_frame = value._owner_frame
+ elif isinstance(value, ElementHandle):
+ handle = value._handle
+ actual_owner = handle._page if handle is not None else value._locator._page
+ owner_frame = value.owner_frame()
+ elif isinstance(value, (list, tuple)):
+ for item in value:
+ validate(item)
+ return
+ elif isinstance(value, dict):
+ for item in value.values():
+ validate(item)
+ return
+ else:
+ return
+
+ if actual_owner is not expected_owner:
+ raise Error(f"{method}: JSHandles can be evaluated only in the context they were created!")
+ if owner_frame is not None and not _same_evaluate_frame(owner_frame, expected_frame):
+ raise Error(f"{method}: JSHandles can be evaluated only in the context they were created!")
+
+ validate(arg)
+
+
def _prepare_evaluate_argument(page: "Page", arg: Any) -> _PreparedEvaluateArgument:
handles: list[JSHandle] = []
temporaries: list[JSHandle] = []
@@ -26620,6 +27016,8 @@ def url(self) -> str:
def evaluate(self, expression: str, arg: Any = None) -> Any:
if self._core is None:
raise Error("worker is not attached")
+ if arg is not None:
+ _ensure_evaluate_argument_context(self, None, arg, method="Worker.evaluate")
prepared = _prepare_evaluate_argument(self, arg) if arg is not None else None
if prepared is not None and prepared.has_handles:
try:
@@ -26641,6 +27039,8 @@ def evaluate(self, expression: str, arg: Any = None) -> Any:
def evaluate_handle(self, expression: str, arg: Any = None) -> JSHandle:
if self._core is None:
raise Error("worker is not attached")
+ if arg is not None:
+ _ensure_evaluate_argument_context(self, None, arg, method="Worker.evaluate_handle")
prepared = _prepare_evaluate_argument(self, arg) if arg is not None else None
if prepared is not None and prepared.has_handles:
try:
From bb5246d5ed134c0ee41b94ea29598bbd3094a082 Mon Sep 17 00:00:00 2001
From: suchintan <3853670+suchintan@users.noreply.github.com>
Date: Sat, 15 Aug 2026 18:45:10 +0000
Subject: [PATCH 3/7] =?UTF-8?q?=F0=9F=94=84=20synced=20local=20'src/'=20wi?=
=?UTF-8?q?th=20remote=20'src/'?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
https://github.com/Skyvern-AI/rustwright-cloud/pull/209
---
src/lib.rs | 2473 +++++++++++++++++++++++++++++++++++++++++++++++-----
1 file changed, 2273 insertions(+), 200 deletions(-)
diff --git a/src/lib.rs b/src/lib.rs
index 77dcf80..d0af95c 100644
--- a/src/lib.rs
+++ b/src/lib.rs
@@ -3161,7 +3161,7 @@ mod tests {
}
impl InputCdpPeer {
- async fn next_command(&mut self) -> Value {
+ async fn next_protocol_command(&mut self) -> Value {
let outgoing = match self.write_rx.recv().await {
Some(outgoing) => outgoing,
None if self.allow_close => {
@@ -3179,6 +3179,15 @@ mod tests {
CdpOutgoing::Close => panic!("unexpected transport close"),
};
let command = serde_json::from_str::(&payload).expect("valid CDP command");
+ let evaluation_source = match command["method"].as_str() {
+ Some("Runtime.evaluate") => command["params"]["expression"].as_str().unwrap_or(""),
+ Some("Runtime.callFunctionOn") => command["params"]["functionDeclaration"]
+ .as_str()
+ .unwrap_or(""),
+ _ => "",
+ };
+ let installing_serializer =
+ evaluation_source.contains("__rustwright_serializer_factory__");
let delete_raw_down = command["method"] == "Input.dispatchKeyEvent"
&& command["params"]["type"] == "rawKeyDown"
&& command["params"]["key"] == "Delete";
@@ -3205,8 +3214,18 @@ mod tests {
json!({ "frameTree": { "frame": { "id": "test-frame" } } })
}
Some("Page.createIsolatedWorld") => json!({ "executionContextId": 1 }),
- Some("Runtime.evaluate") => {
- let expression = command["params"]["expression"].as_str().unwrap_or("");
+ Some("Runtime.evaluate") | Some("Runtime.callFunctionOn")
+ if installing_serializer =>
+ {
+ json!({
+ "result": {
+ "type": "object",
+ "objectId": "input-serializer",
+ }
+ })
+ }
+ Some("Runtime.evaluate") | Some("Runtime.callFunctionOn") => {
+ let expression = evaluation_source;
let value = if expression.contains("doc.activeElement !== fallback") {
Value::Bool(!self.focus_target)
} else if let Some(fill_value) = &self.fill_value {
@@ -3290,7 +3309,9 @@ mod tests {
}
_ => json!({}),
};
- let response_delay = if self.fill_guard_lost
+ let response_delay = if installing_serializer {
+ None
+ } else if self.fill_guard_lost
&& self.delay_resolution_after_guard_loss
&& command["method"] == "Page.getFrameTree"
{
@@ -3301,11 +3322,7 @@ mod tests {
if let Some(delay) = response_delay {
tokio::time::sleep(delay).await;
}
- let response = if command["method"] == "Runtime.evaluate"
- && command["params"]["expression"]
- .as_str()
- .is_some_and(|expression| expression.contains("fallback.focus();"))
- {
+ let response = if evaluation_source.contains("fallback.focus();") {
self.focus_response_error_after_execution
.take()
.map_or_else(
@@ -3315,10 +3332,8 @@ mod tests {
} else {
json!({ "id": command["id"], "result": result })
};
- if let Some(expression) = command["params"]["expression"]
- .as_str()
- .filter(|expression| expression.contains("commitment: 'commencing'"))
- {
+ if evaluation_source.contains("commitment: 'commencing'") {
+ let expression = evaluation_source;
let dispatch_id = action_dispatch_string_literal(expression, "dispatchId");
let token = action_dispatch_string_literal(expression, "token");
let payload = json!({
@@ -3350,6 +3365,26 @@ mod tests {
);
command
}
+ async fn next_command(&mut self) -> Value {
+ loop {
+ let mut command = self.next_protocol_command().await;
+ let source = command["params"]["expression"]
+ .as_str()
+ .or_else(|| command["params"]["functionDeclaration"].as_str())
+ .unwrap_or("");
+ if source.contains("__rustwright_serializer_factory__") {
+ continue;
+ }
+ if command["method"] == "Runtime.callFunctionOn"
+ && source.contains("__rustwright_evaluate_wrapper__")
+ {
+ let source = source.to_string();
+ command["method"] = Value::String("Runtime.evaluate".to_string());
+ command["params"]["expression"] = Value::String(source);
+ }
+ return command;
+ }
+ }
async fn commands(&mut self, count: usize) -> Vec {
let mut commands = Vec::with_capacity(count);
@@ -3377,6 +3412,7 @@ mod tests {
events: events.clone(),
event_log: Arc::clone(&event_log),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -4960,7 +4996,7 @@ multiline-compatible = """4.5.6"""
matches!(
&result,
Err(RwError::Cdp { method, message })
- if method == "Runtime.evaluate" && message == FOCUS_ERROR
+ if method == "Runtime.callFunctionOn" && message == FOCUS_ERROR
),
"{entry_path:?}: {result:?}"
);
@@ -5173,6 +5209,7 @@ multiline-compatible = """4.5.6"""
events: events.clone(),
event_log: Arc::clone(&event_log),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -5242,6 +5279,14 @@ multiline-compatible = """4.5.6"""
}
"Page.createIsolatedWorld" => json!({ "executionContextId": 1 }),
"Runtime.evaluate" => {
+ json!({
+ "result": {
+ "type": "object",
+ "objectId": "typing-serializer",
+ }
+ })
+ }
+ "Runtime.callFunctionOn" => {
json!({ "result": { "type": "boolean", "value": true } })
}
"Input.dispatchKeyEvent" => json!({}),
@@ -5312,6 +5357,7 @@ multiline-compatible = """4.5.6"""
events: events.clone(),
event_log: Arc::clone(&event_log),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -5422,6 +5468,7 @@ multiline-compatible = """4.5.6"""
events: events.clone(),
event_log: Arc::clone(&event_log),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -5820,6 +5867,7 @@ multiline-compatible = """4.5.6"""
events: events.clone(),
event_log: Arc::clone(&event_log),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -5897,6 +5945,7 @@ multiline-compatible = """4.5.6"""
events,
event_log,
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -5997,6 +6046,7 @@ multiline-compatible = """4.5.6"""
events,
event_log: Arc::new(Mutex::new(CdpEventLog::new())),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -7016,6 +7066,20 @@ multiline-compatible = """4.5.6"""
assert!(args.iter().any(|arg| arg == "--use-mock-keychain"));
}
+ #[test]
+ fn chromium_keychain_defaults_honor_user_override_filtering() {
+ let ignored = vec![
+ "--password-store".to_string(),
+ "--use-mock-keychain".to_string(),
+ ];
+
+ assert!(launch_default_arg_ignored(
+ "--password-store=basic",
+ &ignored
+ ));
+ assert!(launch_default_arg_ignored("--use-mock-keychain", &ignored));
+ }
+
#[test]
fn chromium_launch_failure_message_includes_stderr_tail() {
let stderr = NamedTempFile::new().unwrap();
@@ -8667,6 +8731,7 @@ multiline-compatible = """4.5.6"""
events,
event_log: Arc::new(Mutex::new(CdpEventLog::new())),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -8734,6 +8799,7 @@ multiline-compatible = """4.5.6"""
events,
event_log: Arc::new(Mutex::new(CdpEventLog::new())),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -8819,6 +8885,7 @@ multiline-compatible = """4.5.6"""
events: events.clone(),
event_log: Arc::clone(&event_log),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -8896,6 +8963,7 @@ multiline-compatible = """4.5.6"""
events: events.clone(),
event_log: Arc::clone(&event_log),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -8980,6 +9048,7 @@ multiline-compatible = """4.5.6"""
events: events.clone(),
event_log: Arc::clone(&event_log),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -9054,6 +9123,7 @@ multiline-compatible = """4.5.6"""
events,
event_log: Arc::new(Mutex::new(CdpEventLog::new())),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -9228,6 +9298,7 @@ multiline-compatible = """4.5.6"""
events,
event_log: Arc::new(Mutex::new(CdpEventLog::new())),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -9959,11 +10030,10 @@ multiline-compatible = """4.5.6"""
.await
.expect("timed out waiting for CDP test command")
.expect("CDP test command");
- let command: Value = match outgoing {
+ match outgoing {
CdpOutgoing::Text { payload, .. } => serde_json::from_str(&payload).unwrap(),
CdpOutgoing::Close => panic!("unexpected transport close"),
- };
- command
+ }
}
async fn next_command(&mut self, expected_method: &str) -> Value {
@@ -9990,11 +10060,32 @@ multiline-compatible = """4.5.6"""
);
}
+ fn reply_error_with_code(&self, command: &Value, code: i64, message: &str) {
+ dispatch_cdp_payload(
+ json!({
+ "id": command["id"],
+ "error": { "code": code, "message": message },
+ }),
+ Arc::clone(&self.pending),
+ self.events.clone(),
+ Arc::clone(&self.event_log),
+ );
+ }
+
async fn reply_next(&mut self, expected_method: &str, result: Value) {
let command = self.next_command(expected_method).await;
self.reply(&command, result);
}
+ fn observe_runtime_event(&self, event: &Value) {
+ self.page
+ .browser
+ .client
+ .runtime_state
+ .lock()
+ .unwrap()
+ .observe_event(event);
+ }
async fn reply_action_dispatch_binding(&mut self, expected_session_id: &str) {
let binding = self.next_command("Runtime.addBinding").await;
assert_eq!(binding["sessionId"], expected_session_id);
@@ -10004,8 +10095,8 @@ multiline-compatible = """4.5.6"""
);
self.reply(&binding, json!({}));
}
-
fn emit(&self, event: Value) {
+ self.observe_runtime_event(&event);
dispatch_cdp_payload(
event,
Arc::clone(&self.pending),
@@ -10032,6 +10123,7 @@ multiline-compatible = """4.5.6"""
events: events.clone(),
event_log: Arc::clone(&event_log),
traffic_log: Arc::new(Mutex::new(CdpTrafficLog::new())),
+ runtime_state: Arc::new(Mutex::new(CdpRuntimeState::new(None))),
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -11723,6 +11815,31 @@ multiline-compatible = """4.5.6"""
);
}
}
+ "Runtime.callFunctionOn"
+ if session_id == Some("click-child-session") =>
+ {
+ let function = command["params"]["functionDeclaration"]
+ .as_str()
+ .expect("function declaration");
+ if function.contains("__rustwright_serializer_factory__") {
+ harness.reply(
+ &command,
+ json!({
+ "result": {
+ "type": "object",
+ "objectId": "frame-serializer",
+ }
+ }),
+ );
+ } else if function.contains("receives_events") {
+ harness.reply(
+ &command,
+ json!({ "result": { "type": "object", "value": actionable } }),
+ );
+ } else {
+ panic!("unexpected frame function: {function}");
+ }
+ }
"DOM.getBoxModel" => harness.reply(
&command,
json!({ "model": { "border": [10, 15, 14, 15, 14, 21, 10, 21] } }),
@@ -13975,13 +14092,29 @@ multiline-compatible = """4.5.6"""
futures_util::poll!(&mut operation),
std::task::Poll::Pending
));
- let snapshot = harness.next_command("Runtime.evaluate").await;
- assert_eq!(
- snapshot
- .pointer("/params/expression")
- .and_then(Value::as_str),
- Some("document.body.textContent")
+ let serializer = harness.next_command("Runtime.evaluate").await;
+ assert!(serializer
+ .pointer("/params/expression")
+ .and_then(Value::as_str)
+ .is_some_and(|expression| expression.contains("__rustwright_serializer_factory__")));
+ harness.reply(
+ &serializer,
+ json!({
+ "result": {
+ "type": "object",
+ "objectId": "goto-snapshot-serializer",
+ }
+ }),
);
+ assert!(matches!(
+ futures_util::poll!(&mut operation),
+ std::task::Poll::Pending
+ ));
+ let snapshot = harness.next_command("Runtime.callFunctionOn").await;
+ assert!(snapshot
+ .pointer("/params/functionDeclaration")
+ .and_then(Value::as_str)
+ .is_some_and(|function| function.contains("document.body.textContent")));
harness.reply(
&snapshot,
json!({ "result": { "type": "string", "value": "Goto complete" } }),
@@ -14145,13 +14278,29 @@ multiline-compatible = """4.5.6"""
futures_util::poll!(&mut operation),
std::task::Poll::Pending
));
- let snapshot = harness.next_command("Runtime.evaluate").await;
- assert_eq!(
- snapshot
- .pointer("/params/expression")
- .and_then(Value::as_str),
- Some("document.body.textContent")
+ let serializer = harness.next_command("Runtime.evaluate").await;
+ assert!(serializer
+ .pointer("/params/expression")
+ .and_then(Value::as_str)
+ .is_some_and(|expression| expression.contains("__rustwright_serializer_factory__")));
+ harness.reply(
+ &serializer,
+ json!({
+ "result": {
+ "type": "object",
+ "objectId": "reload-snapshot-serializer",
+ }
+ }),
);
+ assert!(matches!(
+ futures_util::poll!(&mut operation),
+ std::task::Poll::Pending
+ ));
+ let snapshot = harness.next_command("Runtime.callFunctionOn").await;
+ assert!(snapshot
+ .pointer("/params/functionDeclaration")
+ .and_then(Value::as_str)
+ .is_some_and(|function| function.contains("document.body.textContent")));
harness.reply(
&snapshot,
json!({ "result": { "type": "string", "value": "Status waiting" } }),
@@ -14879,6 +15028,17 @@ multiline-compatible = """4.5.6"""
harness
.reply_next(
"Runtime.evaluate",
+ json!({
+ "result": {
+ "type": "object",
+ "objectId": "dialog-action-serializer",
+ }
+ }),
+ )
+ .await;
+ harness
+ .reply_next(
+ "Runtime.callFunctionOn",
json!({
"result": {
"type": "object",
@@ -15057,9 +15217,20 @@ multiline-compatible = """4.5.6"""
Some("frame-main")
);
harness.reply(&utility_world, json!({ "executionContextId": 1 }));
+ harness
+ .reply_next(
+ "Runtime.evaluate",
+ json!({
+ "result": {
+ "type": "object",
+ "objectId": "public-action-serializer",
+ }
+ }),
+ )
+ .await;
if matches!(action, PublicSettledAction::Scroll) {
- let dispatch = harness.next_command("Runtime.evaluate").await;
+ let dispatch = harness.next_command("Runtime.callFunctionOn").await;
assert!(
harness.events.receiver_count() >= 1,
"scroll settlement must subscribe before its dispatch"
@@ -15079,7 +15250,7 @@ multiline-compatible = """4.5.6"""
harness
.reply_next(
- "Runtime.evaluate",
+ "Runtime.callFunctionOn",
json!({
"result": {
"type": "object",
@@ -15135,7 +15306,7 @@ multiline-compatible = """4.5.6"""
harness.reply(&utility_world, json!({ "executionContextId": 3 }));
harness
.reply_next(
- "Runtime.evaluate",
+ "Runtime.callFunctionOn",
json!({
"result": {
"type": "string",
@@ -15180,7 +15351,7 @@ multiline-compatible = """4.5.6"""
harness.reply(&utility_world, json!({ "executionContextId": 4 }));
harness
.reply_next(
- "Runtime.evaluate",
+ "Runtime.callFunctionOn",
json!({ "result": { "type": "boolean", "value": true } }),
)
.await;
@@ -15213,6 +15384,597 @@ multiline-compatible = """4.5.6"""
}
}
+ #[tokio::test]
+ async fn evaluate_data_result_is_one_runtime_command_and_handle_path_is_unchanged() {
+ let mut harness = navigation_test_harness(4);
+
+ let client = Arc::clone(&harness.page.browser.client);
+ let warmup = tokio::spawn(async move {
+ evaluate_expression_in_session_before(
+ &client,
+ "page-session",
+ None,
+ make_evaluate_expression("0", None),
+ OperationDeadline::new(Duration::from_secs(1)),
+ )
+ .await
+ });
+ let install_command = harness.next_command("Runtime.evaluate").await;
+ assert_eq!(install_command["params"]["returnByValue"], false);
+ assert_eq!(
+ install_command["params"]["expression"],
+ format!("({})()", runtime_value_serializer_factory())
+ );
+ harness.reply(
+ &install_command,
+ json!({
+ "result": {
+ "type": "function",
+ "objectId": "serializer-main",
+ },
+ }),
+ );
+ let warmup_command = harness.next_command("Runtime.callFunctionOn").await;
+ assert_eq!(warmup_command["params"]["objectId"], "serializer-main");
+ assert_eq!(
+ warmup_command["params"]["arguments"],
+ json!([{ "objectId": "serializer-main" }])
+ );
+ harness.reply(
+ &warmup_command,
+ json!({ "result": { "type": "number", "value": 0 } }),
+ );
+ assert_eq!(warmup.await.unwrap().unwrap(), "0");
+
+ let client = Arc::clone(&harness.page.browser.client);
+ let object_evaluate = tokio::spawn(async move {
+ evaluate_expression_in_session_before(
+ &client,
+ "page-session",
+ None,
+ "({ answer: 42 })".to_string(),
+ OperationDeadline::new(Duration::from_secs(1)),
+ )
+ .await
+ });
+ let object_command = harness.next_command("Runtime.callFunctionOn").await;
+ assert_eq!(object_command["params"]["objectId"], "serializer-main");
+ assert_eq!(object_command["params"]["returnByValue"], true);
+ assert_eq!(
+ object_command["params"]["arguments"],
+ json!([{ "objectId": "serializer-main" }])
+ );
+ assert_eq!(
+ object_command["params"]["functionDeclaration"],
+ serialize_evaluate_result_function("({ answer: 42 })")
+ );
+ assert!(object_command["params"]["functionDeclaration"]
+ .as_str()
+ .unwrap()
+ .contains("__rustwright_evaluate_wrapper__"));
+ assert!(!object_command["params"]["functionDeclaration"]
+ .as_str()
+ .unwrap()
+ .contains(RUNTIME_VALUE_SERIALIZER));
+ harness.reply(
+ &object_command,
+ json!({
+ "result": {
+ "type": "object",
+ "value": {
+ "__rustwright_cdp_object__": 1,
+ "entries": { "answer": 42 },
+ },
+ },
+ }),
+ );
+ assert_eq!(
+ tokio::time::timeout(Duration::from_secs(1), object_evaluate)
+ .await
+ .expect("steady object evaluate should finish after one Runtime reply")
+ .unwrap()
+ .unwrap(),
+ json!({
+ "__rustwright_cdp_object__": 1,
+ "entries": { "answer": 42 },
+ })
+ .to_string()
+ );
+ assert!(matches!(
+ harness.write_rx.try_recv(),
+ Err(tokio::sync::mpsc::error::TryRecvError::Empty)
+ ));
+
+ let client = Arc::clone(&harness.page.browser.client);
+ let primitive_evaluate = tokio::spawn(async move {
+ evaluate_expression_in_session_before(
+ &client,
+ "page-session",
+ None,
+ make_evaluate_expression("1", None),
+ OperationDeadline::new(Duration::from_secs(1)),
+ )
+ .await
+ });
+ let primitive_command = harness.next_command("Runtime.callFunctionOn").await;
+ assert_eq!(primitive_command["params"]["objectId"], "serializer-main");
+ assert_eq!(
+ primitive_command["params"]["arguments"],
+ json!([{ "objectId": "serializer-main" }])
+ );
+ harness.reply(
+ &primitive_command,
+ json!({ "result": { "type": "number", "value": 1 } }),
+ );
+ assert_eq!(
+ tokio::time::timeout(Duration::from_secs(1), primitive_evaluate)
+ .await
+ .expect("primitive evaluate should keep its one-command path")
+ .unwrap()
+ .unwrap(),
+ "1"
+ );
+ assert!(matches!(
+ harness.write_rx.try_recv(),
+ Err(tokio::sync::mpsc::error::TryRecvError::Empty)
+ ));
+
+ let page = Arc::clone(&harness.page);
+ let call_function_evaluate = tokio::spawn(async move {
+ evaluate_locator_for_element_handle(
+ page,
+ "element-handle-1".to_string(),
+ Some("page-session".to_string()),
+ r##"{"kind":"css","selector":"#target"}"##.to_string(),
+ 0,
+ "return { answer: 42 };".to_string(),
+ Duration::from_secs(1),
+ true,
+ false,
+ Duration::ZERO,
+ )
+ .await
+ });
+ let call_function_command = harness.next_command("Runtime.callFunctionOn").await;
+ assert_eq!(
+ call_function_command["params"]["objectId"],
+ "element-handle-1"
+ );
+ assert_eq!(
+ call_function_command["params"]["arguments"],
+ json!([{ "objectId": "serializer-main" }])
+ );
+ assert!(call_function_command["params"]["functionDeclaration"]
+ .as_str()
+ .unwrap()
+ .contains("__rustwright_evaluate_wrapper__"));
+ assert!(!call_function_command["params"]["functionDeclaration"]
+ .as_str()
+ .unwrap()
+ .contains(RUNTIME_VALUE_SERIALIZER));
+ harness.reply(
+ &call_function_command,
+ json!({
+ "result": {
+ "type": "object",
+ "value": {
+ "__rustwright_cdp_object__": 1,
+ "entries": { "answer": 42 },
+ },
+ },
+ }),
+ );
+ assert_eq!(
+ call_function_evaluate.await.unwrap().unwrap(),
+ json!({
+ "__rustwright_cdp_object__": 1,
+ "entries": { "answer": 42 },
+ })
+ .to_string()
+ );
+
+ let client = Arc::clone(&harness.page.browser.client);
+ let handle_evaluate = tokio::spawn(async move {
+ evaluate_handle_expression_in_session(
+ &client,
+ "page-session",
+ "globalThis".to_string(),
+ OperationDeadline::new(Duration::from_secs(1)),
+ )
+ .await
+ });
+ let handle_command = harness.next_command("Runtime.evaluate").await;
+ assert_eq!(handle_command["params"]["returnByValue"], false);
+ assert_eq!(handle_command["params"]["expression"], "globalThis");
+ harness.reply(
+ &handle_command,
+ json!({ "result": { "type": "object", "objectId": "remote-handle-1" } }),
+ );
+ assert_eq!(
+ handle_evaluate.await.unwrap().unwrap(),
+ json!({ "type": "object", "objectId": "remote-handle-1" })
+ );
+ assert!(matches!(
+ harness.write_rx.try_recv(),
+ Err(tokio::sync::mpsc::error::TryRecvError::Empty)
+ ));
+ }
+
+ #[tokio::test]
+ async fn declaration_helper_user_wrapper_name_survives_inline_exception_conversion() {
+ let mut harness = navigation_test_harness(4);
+ let client = Arc::clone(&harness.page.browser.client);
+ let source = concat!(
+ "const marker = 1; function __rustwright_evaluate_wrapper__() {",
+ " throw new Error('user boom');",
+ " } __rustwright_evaluate_wrapper__();"
+ );
+ let expression = make_evaluate_expression(source, None);
+ assert!(is_script_goal_evaluate_expression(&expression));
+ let evaluate = tokio::spawn(async move {
+ evaluate_expression_in_session_before(
+ &client,
+ "page-session",
+ None,
+ expression,
+ OperationDeadline::new(Duration::from_secs(1)),
+ )
+ .await
+ });
+ let command = harness.next_command("Runtime.evaluate").await;
+ let description = concat!(
+ "Error: user boom\n",
+ " at __rustwright_evaluate_wrapper__ (:2:9)\n",
+ " at caller (https://example.test/app.js:7:3)"
+ );
+ harness.reply(
+ &command,
+ json!({
+ "exceptionDetails": {
+ "scriptId": "user-script",
+ "exception": { "description": description },
+ "stackTrace": {
+ "callFrames": [{
+ "functionName": "__rustwright_evaluate_wrapper__",
+ "scriptId": "user-script",
+ "url": "",
+ "lineNumber": 1,
+ "columnNumber": 8,
+ }],
+ },
+ },
+ "result": { "type": "object", "subtype": "error" },
+ }),
+ );
+
+ assert_eq!(
+ evaluate.await.unwrap().unwrap_err().to_string(),
+ description
+ );
+ }
+
+ #[tokio::test]
+ async fn serializer_cache_is_one_install_per_realm_across_new_handles() {
+ let mut harness = navigation_test_harness(4);
+ let client = Arc::clone(&harness.page.browser.client);
+ let realm = client.serializer_realm_key("page-session", Some("frame:main"));
+
+ for index in 0..8 {
+ let client = Arc::clone(&client);
+ let realm = realm.clone();
+ let object_id = format!("handle-{index}");
+ let call = tokio::spawn(async move {
+ call_function_with_serialized_result_before(
+ &client,
+ &realm,
+ &object_id,
+ "function() { return this.value; }",
+ None,
+ OperationDeadline::new(Duration::from_secs(1)),
+ )
+ .await
+ });
+ if index == 0 {
+ let install = harness.next_command("Runtime.callFunctionOn").await;
+ assert_eq!(install["params"]["objectId"], "handle-0");
+ assert!(install["params"]["functionDeclaration"]
+ .as_str()
+ .unwrap()
+ .contains("__rustwright_serializer_factory__"));
+ harness.reply(
+ &install,
+ json!({
+ "result": {
+ "type": "function",
+ "objectId": "realm-serializer",
+ },
+ }),
+ );
+ }
+ let command = harness.next_command("Runtime.callFunctionOn").await;
+ assert_eq!(command["params"]["objectId"], format!("handle-{index}"));
+ assert_eq!(
+ command["params"]["arguments"],
+ json!([{ "objectId": "realm-serializer" }])
+ );
+ harness.reply(
+ &command,
+ json!({ "result": { "type": "number", "value": index } }),
+ );
+ assert_eq!(
+ call.await.unwrap().unwrap()["result"]["value"],
+ json!(index)
+ );
+ }
+
+ assert_eq!(client.runtime_state.lock().unwrap().serializers.len(), 1);
+ assert_eq!(
+ client
+ .runtime_state
+ .lock()
+ .unwrap()
+ .serializer_install_locks
+ .len(),
+ 1
+ );
+ assert!(matches!(
+ harness.write_rx.try_recv(),
+ Err(tokio::sync::mpsc::error::TryRecvError::Empty)
+ ));
+ }
+
+ #[tokio::test]
+ async fn detached_realm_cannot_commit_in_flight_serializer_install() {
+ let mut harness = navigation_test_harness(4);
+ let client = Arc::clone(&harness.page.browser.client);
+
+ for index in 0..8 {
+ let realm =
+ client.serializer_realm_key("page-session", Some(&format!("frame:child-{index}")));
+ let install_client = Arc::clone(&client);
+ let install_realm = realm.clone();
+ let install = tokio::spawn(async move {
+ serializer_for_realm(
+ &install_client,
+ &install_realm,
+ None,
+ None,
+ OperationDeadline::new(Duration::from_secs(1)),
+ )
+ .await
+ });
+ let install_command = harness.next_command("Runtime.evaluate").await;
+ harness.reply(
+ &install_command,
+ json!({
+ "result": {
+ "type": "function",
+ "objectId": format!("invalidated-serializer-{index}"),
+ },
+ }),
+ );
+ // The current-thread test runtime does not resume the install task until
+ // this test yields, so the detach deterministically wins the commit.
+ harness.observe_runtime_event(&json!({
+ "method": "Target.detachedFromTarget",
+ "params": { "sessionId": "page-session" },
+ }));
+ {
+ let state = client.runtime_state.lock().unwrap();
+ assert!(!state.serializer_generations.contains_key(&realm));
+ assert!(state.serializer_install_locks.is_empty());
+ }
+
+ assert_eq!(
+ install.await.unwrap().unwrap_err().to_string(),
+ "Execution context was destroyed, most likely because of a navigation."
+ );
+ let state = client.runtime_state.lock().unwrap();
+ assert!(state.serializers.is_empty());
+ assert!(
+ state.serializer_install_locks.is_empty(),
+ "cycle {index} retained an install lock"
+ );
+ assert!(
+ state.serializer_generations.is_empty(),
+ "cycle {index} retained a serializer generation"
+ );
+ }
+ }
+
+ #[tokio::test]
+ async fn failed_serializer_install_removes_exact_install_lock() {
+ let mut harness = navigation_test_harness(4);
+ let client = Arc::clone(&harness.page.browser.client);
+ let realm = client.serializer_realm_key("page-session", Some("frame:main"));
+ let install_client = Arc::clone(&client);
+ let install_realm = realm.clone();
+ let install = tokio::spawn(async move {
+ serializer_for_realm(
+ &install_client,
+ &install_realm,
+ None,
+ None,
+ OperationDeadline::new(Duration::from_secs(1)),
+ )
+ .await
+ });
+ let install_command = harness.next_command("Runtime.evaluate").await;
+ harness.reply_error_with_code(&install_command, -32_000, "serializer install failed");
+
+ assert!(install.await.unwrap().is_err());
+ let state = client.runtime_state.lock().unwrap();
+ assert!(state.serializer_install_locks.is_empty());
+ assert!(state.serializer_generations.is_empty());
+ }
+
+ #[tokio::test]
+ async fn serializer_navigation_eviction_releases_live_object() {
+ let mut harness = navigation_test_harness(4);
+ let client = Arc::clone(&harness.page.browser.client);
+ let (release_tx, release_rx) = mpsc::unbounded_channel();
+ client.runtime_state.lock().unwrap().release_tx = Some(release_tx);
+ spawn_serializer_release_pump(release_rx, client.write_tx.clone());
+ let realm = client.serializer_realm_key("page-session", Some("frame:main"));
+ let (generation, install_lock) = client.serializer_install_lock(&realm);
+ {
+ let _install_guard = install_lock.lock().await;
+ assert!(client.remember_serializer(
+ &realm,
+ generation,
+ &install_lock,
+ "live-serializer".to_string(),
+ Some("7".to_string()),
+ ));
+ }
+
+ harness.observe_runtime_event(&json!({
+ "method": "Page.frameNavigated",
+ "sessionId": "page-session",
+ "params": { "frame": { "id": "main", "loaderId": "loader-a" } }
+ }));
+ harness.observe_runtime_event(&json!({
+ "method": "Page.frameNavigated",
+ "sessionId": "page-session",
+ "params": { "frame": { "id": "main", "loaderId": "loader-b" } }
+ }));
+
+ let release = harness.next_command("Runtime.releaseObject").await;
+ assert_eq!(release["sessionId"], "page-session");
+ assert_eq!(release["params"]["objectId"], "live-serializer");
+ assert!(client.serializer_handle(&realm).is_none());
+ assert!(client
+ .runtime_state
+ .lock()
+ .unwrap()
+ .serializer_install_locks
+ .is_empty());
+ }
+
+ #[tokio::test]
+ async fn execution_context_destroyed_never_retries_user_code() {
+ let mut harness = navigation_test_harness(4);
+ let client = Arc::clone(&harness.page.browser.client);
+ let evaluate_client = Arc::clone(&client);
+ let evaluate = tokio::spawn(async move {
+ evaluate_expression_in_session_before(
+ &evaluate_client,
+ "page-session",
+ Some("frame:main"),
+ make_evaluate_expression("globalThis.beaconCount += 1", None),
+ OperationDeadline::new(Duration::from_secs(1)),
+ )
+ .await
+ });
+ let install = harness.next_command("Runtime.evaluate").await;
+ harness.reply(
+ &install,
+ json!({
+ "result": {
+ "type": "function",
+ "objectId": "serializer-before-navigation",
+ },
+ }),
+ );
+ let user_command = harness.next_command("Runtime.callFunctionOn").await;
+ harness.reply_error_with_code(&user_command, -32_000, "Execution context was destroyed.");
+
+ assert_eq!(
+ evaluate.await.unwrap().unwrap_err().to_string(),
+ "Execution context was destroyed, most likely because of a navigation."
+ );
+ assert!(matches!(
+ harness.write_rx.try_recv(),
+ Err(tokio::sync::mpsc::error::TryRecvError::Empty)
+ ));
+ assert!(client
+ .serializer_handle(&client.serializer_realm_key("page-session", Some("frame:main"),))
+ .is_none());
+ }
+
+ #[tokio::test]
+ async fn evaluate_serializer_cache_reinstalls_once_after_navigation_recreates_context() {
+ let mut harness = navigation_test_harness(4);
+ let client = Arc::clone(&harness.page.browser.client);
+ let initial = tokio::spawn(async move {
+ evaluate_expression_in_session_before(
+ &client,
+ "page-session",
+ None,
+ make_evaluate_expression("1", None),
+ OperationDeadline::new(Duration::from_secs(1)),
+ )
+ .await
+ });
+ let install_command = harness.next_command("Runtime.evaluate").await;
+ harness.reply(
+ &install_command,
+ json!({
+ "result": {
+ "type": "function",
+ "objectId": "serializer-before-navigation",
+ },
+ }),
+ );
+ let initial_call = harness.next_command("Runtime.callFunctionOn").await;
+ assert_eq!(
+ initial_call["params"]["objectId"],
+ "serializer-before-navigation"
+ );
+ harness.reply(
+ &initial_call,
+ json!({ "result": { "type": "number", "value": 1 } }),
+ );
+ assert_eq!(initial.await.unwrap().unwrap(), "1");
+
+ let client = Arc::clone(&harness.page.browser.client);
+ let after_navigation = tokio::spawn(async move {
+ evaluate_expression_in_session_before(
+ &client,
+ "page-session",
+ None,
+ make_evaluate_expression("2", None),
+ OperationDeadline::new(Duration::from_secs(1)),
+ )
+ .await
+ });
+ let stale_cache_command = harness.next_command("Runtime.callFunctionOn").await;
+ assert_eq!(
+ stale_cache_command["params"]["objectId"],
+ "serializer-before-navigation"
+ );
+ harness.reply_error_with_code(
+ &stale_cache_command,
+ -32_000,
+ "Could not find object with given id",
+ );
+ let reinstall_command = harness.next_command("Runtime.evaluate").await;
+ assert_eq!(reinstall_command["params"]["returnByValue"], false);
+ harness.reply(
+ &reinstall_command,
+ json!({
+ "result": {
+ "type": "function",
+ "objectId": "serializer-after-navigation",
+ },
+ }),
+ );
+ let retry_command = harness.next_command("Runtime.callFunctionOn").await;
+ assert_eq!(
+ retry_command["params"]["objectId"],
+ "serializer-after-navigation"
+ );
+ harness.reply(
+ &retry_command,
+ json!({ "result": { "type": "number", "value": 2 } }),
+ );
+ assert_eq!(after_navigation.await.unwrap().unwrap(), "2");
+ assert!(matches!(
+ harness.write_rx.try_recv(),
+ Err(tokio::sync::mpsc::error::TryRecvError::Empty)
+ ));
+ }
+
#[tokio::test]
async fn navigation_wait_completes_on_expected_same_document_url() {
let harness = navigation_test_harness(4);
@@ -17008,6 +17770,378 @@ impl CdpTrafficLog {
}
}
+#[derive(Clone, Debug, Eq, Hash, PartialEq)]
+struct SerializerRealmKey {
+ session_id: String,
+ realm_identity: String,
+}
+
+impl SerializerRealmKey {
+ fn new(session_id: &str, realm_identity: impl Into) -> Self {
+ Self {
+ session_id: session_id.to_string(),
+ realm_identity: realm_identity.into(),
+ }
+ }
+
+ fn in_world(mut self, world_name: &str) -> Self {
+ self.realm_identity = format!("{}|world:{world_name}", self.realm_identity);
+ self
+ }
+
+ fn frame_id(&self) -> Option<&str> {
+ self.realm_identity
+ .strip_prefix("frame:")
+ .and_then(|identity| identity.split("|world:").next())
+ }
+}
+
+#[cfg(test)]
+mod serializer_realm_key_tests {
+ use super::*;
+
+ #[test]
+ fn utility_world_key_does_not_alias_the_default_frame_realm() {
+ let default_realm = SerializerRealmKey::new("page-session", "frame:main");
+ let utility_realm = default_realm.clone().in_world(FRAME_UTILITY_WORLD_NAME);
+
+ assert_ne!(utility_realm, default_realm);
+ assert_eq!(default_realm.frame_id(), Some("main"));
+ assert_eq!(utility_realm.frame_id(), Some("main"));
+ }
+}
+
+#[derive(Clone, Debug)]
+struct SerializerCacheEntry {
+ object_id: String,
+ execution_context_id: Option,
+}
+
+#[derive(Clone, Debug)]
+struct RuntimeExecutionRealm {
+ realm_identity: String,
+ frame_id: Option,
+}
+
+#[derive(Clone, Debug)]
+struct SerializerRelease {
+ session_id: String,
+ object_id: String,
+}
+
+struct CdpRuntimeState {
+ serializers: HashMap,
+ serializer_install_locks: HashMap>>,
+ serializer_generations: HashMap,
+ next_serializer_generation: u64,
+ execution_realms: HashMap<(String, String), RuntimeExecutionRealm>,
+ session_realms: HashMap,
+ frame_loaders: HashMap<(String, String), String>,
+ release_tx: Option>,
+ #[cfg(any(test, feature = "test-support"))]
+ serializer_release_count: usize,
+}
+
+impl CdpRuntimeState {
+ fn new(release_tx: Option>) -> Self {
+ Self {
+ serializers: HashMap::new(),
+ serializer_install_locks: HashMap::new(),
+ serializer_generations: HashMap::new(),
+ next_serializer_generation: 0,
+ execution_realms: HashMap::new(),
+ session_realms: HashMap::new(),
+ frame_loaders: HashMap::new(),
+ #[cfg(any(test, feature = "test-support"))]
+ serializer_release_count: 0,
+ release_tx,
+ }
+ }
+
+ fn default_realm_identity(&self, session_id: &str) -> Option {
+ self.session_realms.get(session_id).cloned()
+ }
+
+ fn execution_context_for_realm(&self, realm: &SerializerRealmKey) -> Option {
+ self.execution_realms
+ .iter()
+ .find_map(|((session_id, context_id), execution_realm)| {
+ (session_id == &realm.session_id
+ && execution_realm.realm_identity == realm.realm_identity)
+ .then(|| context_id.clone())
+ })
+ }
+
+ fn frame_id_for_execution_context(
+ &self,
+ session_id: &str,
+ execution_context_id: &str,
+ ) -> Option {
+ self.execution_realms
+ .get(&(session_id.to_string(), execution_context_id.to_string()))
+ .and_then(|realm| realm.frame_id.clone())
+ }
+
+ fn next_serializer_generation(&mut self) -> u64 {
+ self.next_serializer_generation = self
+ .next_serializer_generation
+ .checked_add(1)
+ .expect("serializer generation counter exhausted");
+ self.next_serializer_generation
+ }
+
+ fn ensure_serializer_generation(&mut self, realm: &SerializerRealmKey) -> u64 {
+ if let Some(generation) = self.serializer_generations.get(realm) {
+ return *generation;
+ }
+ let generation = self.next_serializer_generation();
+ self.serializer_generations
+ .insert(realm.clone(), generation);
+ generation
+ }
+
+ fn enqueue_serializer_release(&mut self, release: SerializerRelease) {
+ let release_enqueued = self
+ .release_tx
+ .as_ref()
+ .is_some_and(|release_tx| release_tx.send(release).is_ok());
+ #[cfg(any(test, feature = "test-support"))]
+ {
+ self.serializer_release_count += usize::from(release_enqueued);
+ }
+ #[cfg(not(any(test, feature = "test-support")))]
+ let _ = release_enqueued;
+ }
+
+ fn evict_serializer_realms(
+ &mut self,
+ realms: impl IntoIterator- ,
+ release_live: bool,
+ ) {
+ for realm in realms.into_iter().collect::>() {
+ let had_generation = self.serializer_generations.remove(&realm).is_some();
+ let had_lock = self.serializer_install_locks.remove(&realm).is_some();
+ let entry = self.serializers.remove(&realm);
+ if !had_generation && !had_lock && entry.is_none() {
+ continue;
+ }
+ if release_live {
+ if let Some(entry) = entry {
+ self.enqueue_serializer_release(SerializerRelease {
+ session_id: realm.session_id,
+ object_id: entry.object_id,
+ });
+ }
+ }
+ }
+ }
+
+ fn evict_matching_serializer_realms(
+ &mut self,
+ mut should_evict: impl FnMut(&SerializerRealmKey) -> bool,
+ release_live: bool,
+ ) {
+ let evicted = self
+ .serializers
+ .keys()
+ .chain(self.serializer_install_locks.keys())
+ .chain(self.serializer_generations.keys())
+ .filter(|realm| should_evict(realm))
+ .cloned()
+ .collect::>();
+ self.evict_serializer_realms(evicted, release_live);
+ }
+
+ fn observe_event(&mut self, event: &Value) {
+ let method = event.get("method").and_then(Value::as_str).unwrap_or("");
+ if method == "Target.attachedToTarget" {
+ let child_session_id = event.pointer("/params/sessionId").and_then(Value::as_str);
+ let target_info = event.pointer("/params/targetInfo");
+ if let (Some(child_session_id), Some(target_info)) = (child_session_id, target_info) {
+ let target_id = target_info.get("targetId").and_then(Value::as_str);
+ let target_type = target_info.get("type").and_then(Value::as_str);
+ if let (Some(target_id), Some(target_type)) = (target_id, target_type) {
+ let realm_identity = match target_type {
+ "page" | "iframe" => Some(format!("frame:{target_id}")),
+ "worker" | "service_worker" | "shared_worker" => {
+ Some(format!("worker:{target_id}"))
+ }
+ _ => None,
+ };
+ if let Some(realm_identity) = realm_identity {
+ self.session_realms
+ .insert(child_session_id.to_string(), realm_identity);
+ }
+ }
+ }
+ return;
+ }
+
+ if method == "Target.detachedFromTarget" {
+ if let Some(detached_session_id) =
+ event.pointer("/params/sessionId").and_then(Value::as_str)
+ {
+ self.evict_session(detached_session_id, false);
+ }
+ return;
+ }
+
+ let Some(session_id) = event.get("sessionId").and_then(Value::as_str) else {
+ return;
+ };
+ let params = event.get("params").unwrap_or(&Value::Null);
+ match method {
+ "Runtime.executionContextCreated" => {
+ let context = params.get("context").unwrap_or(&Value::Null);
+ let Some(context_id) = context.get("id") else {
+ return;
+ };
+ let context_id = context_id.to_string();
+ let frame_id = context
+ .pointer("/auxData/frameId")
+ .and_then(Value::as_str)
+ .map(ToString::to_string);
+ let realm_identity = frame_id
+ .as_deref()
+ .map(|frame_id| format!("frame:{frame_id}"))
+ .or_else(|| self.session_realms.get(session_id).cloned())
+ .unwrap_or_else(|| format!("session:{session_id}"));
+ self.execution_realms.insert(
+ (session_id.to_string(), context_id),
+ RuntimeExecutionRealm {
+ realm_identity,
+ frame_id,
+ },
+ );
+ }
+ "Runtime.executionContextDestroyed" => {
+ let context_id = params
+ .get("executionContextId")
+ .map(Value::to_string)
+ .or_else(|| {
+ params
+ .get("executionContextUniqueId")
+ .and_then(Value::as_str)
+ .map(ToString::to_string)
+ });
+ if let Some(context_id) = context_id {
+ let execution_realm = self
+ .execution_realms
+ .remove(&(session_id.to_string(), context_id.clone()));
+ let mut evicted = self
+ .serializers
+ .iter()
+ .filter(|(realm, entry)| {
+ realm.session_id == session_id
+ && entry.execution_context_id.as_deref()
+ == Some(context_id.as_str())
+ })
+ .map(|(realm, _)| realm.clone())
+ .collect::>();
+ if let Some(execution_realm) = execution_realm {
+ let realm =
+ SerializerRealmKey::new(session_id, execution_realm.realm_identity);
+ if self.serializer_generations.contains_key(&realm)
+ || self.serializer_install_locks.contains_key(&realm)
+ || self.serializers.contains_key(&realm)
+ {
+ evicted.push(realm);
+ }
+ }
+ self.evict_serializer_realms(evicted, false);
+ }
+ }
+ "Runtime.executionContextsCleared" => self.evict_session(session_id, false),
+ "Page.frameNavigated" => {
+ let frame = params.get("frame").unwrap_or(&Value::Null);
+ let frame_id = frame.get("id").and_then(Value::as_str);
+ let loader_id = frame.get("loaderId").and_then(Value::as_str);
+ if let (Some(frame_id), Some(loader_id)) = (frame_id, loader_id) {
+ let loader_key = (session_id.to_string(), frame_id.to_string());
+ let loader_changed = self
+ .frame_loaders
+ .insert(loader_key, loader_id.to_string())
+ .is_some_and(|previous| previous != loader_id);
+ if loader_changed {
+ self.evict_frame(session_id, frame_id, true);
+ }
+ }
+ }
+ "Page.frameDetached" => {
+ let reason = params.get("reason").and_then(Value::as_str);
+ let release_live = match reason {
+ Some("remove") => Some(false),
+ Some("swap") => Some(true),
+ _ => None,
+ };
+ if let (Some(frame_id), Some(release_live)) =
+ (params.get("frameId").and_then(Value::as_str), release_live)
+ {
+ self.evict_frame(session_id, frame_id, release_live);
+ if reason == Some("remove") {
+ self.frame_loaders
+ .remove(&(session_id.to_string(), frame_id.to_string()));
+ }
+ }
+ }
+ "Page.frameSwapped" | "Page.frameSwappedByActivation" => {
+ if let Some(frame_id) = params.get("frameId").and_then(Value::as_str) {
+ self.evict_frame(session_id, frame_id, true);
+ }
+ }
+ _ => {}
+ }
+ }
+
+ fn evict_frame(&mut self, session_id: &str, frame_id: &str, release_live: bool) {
+ self.evict_matching_serializer_realms(
+ |realm| realm.session_id == session_id && realm.frame_id() == Some(frame_id),
+ release_live,
+ );
+ self.execution_realms
+ .retain(|(realm_session_id, _), realm| {
+ realm_session_id != session_id || realm.frame_id.as_deref() != Some(frame_id)
+ });
+ }
+
+ fn evict_session(&mut self, session_id: &str, release_live: bool) {
+ self.evict_matching_serializer_realms(|realm| realm.session_id == session_id, release_live);
+ self.execution_realms
+ .retain(|(realm_session_id, _), _| realm_session_id != session_id);
+ self.session_realms.remove(session_id);
+ self.frame_loaders
+ .retain(|(loader_session_id, _), _| loader_session_id != session_id);
+ }
+}
+
+fn spawn_serializer_release_pump(
+ mut release_rx: mpsc::UnboundedReceiver,
+ write_tx: mpsc::UnboundedSender,
+) {
+ static NEXT_RELEASE_ID: AtomicU64 = AtomicU64::new(8_000_000_000_000_000);
+ tokio::spawn(async move {
+ while let Some(release) = release_rx.recv().await {
+ let id = NEXT_RELEASE_ID.fetch_add(1, Ordering::SeqCst);
+ let payload = json!({
+ "id": id,
+ "method": "Runtime.releaseObject",
+ "params": { "objectId": release.object_id },
+ "sessionId": release.session_id,
+ });
+ if write_tx
+ .send(CdpOutgoing::Text {
+ payload: payload.to_string(),
+ tracker: None,
+ diagnostic_id: None,
+ })
+ .is_err()
+ {
+ break;
+ }
+ }
+ });
+}
+
struct CdpClientDiagnosticSnapshot {
captured_at: Instant,
traffic: Vec,
@@ -17127,6 +18261,7 @@ struct CdpClient {
events: broadcast::Sender,
event_log: Arc>,
traffic_log: Arc>,
+ runtime_state: Arc>,
next_id: AtomicU64,
sent_runtime_enable_count: AtomicU64,
sent_target_close_count: AtomicU64,
@@ -17374,6 +18509,7 @@ fn dispatch_cdp_payload_with_diagnostics(
events: broadcast::Sender,
event_log: Arc>,
traffic_log: Arc>,
+ runtime_state: Arc>,
) {
if let Some(id) = payload.get("id").and_then(Value::as_u64) {
let sender = pending.lock().unwrap().remove(&id);
@@ -17402,6 +18538,7 @@ fn dispatch_cdp_payload_with_diagnostics(
let _ = sender.send(result);
}
} else {
+ runtime_state.lock().unwrap().observe_event(&payload);
traffic_log
.lock()
.unwrap()
@@ -17426,6 +18563,7 @@ fn dispatch_cdp_payload(
events,
event_log,
Arc::new(Mutex::new(CdpTrafficLog::new())),
+ Arc::new(Mutex::new(CdpRuntimeState::new(None))),
);
}
@@ -17486,6 +18624,140 @@ fn ensure_ws_request_path(
}
impl CdpClient {
+ fn serializer_realm_key(
+ &self,
+ session_id: &str,
+ realm_identity: Option<&str>,
+ ) -> SerializerRealmKey {
+ let realm_identity = realm_identity
+ .map(ToString::to_string)
+ .or_else(|| {
+ self.runtime_state
+ .lock()
+ .unwrap()
+ .default_realm_identity(session_id)
+ })
+ .unwrap_or_else(|| format!("session:{session_id}"));
+ SerializerRealmKey::new(session_id, realm_identity)
+ }
+
+ fn serializer_handle(&self, realm: &SerializerRealmKey) -> Option {
+ self.runtime_state
+ .lock()
+ .unwrap()
+ .serializers
+ .get(realm)
+ .map(|entry| entry.object_id.clone())
+ }
+
+ fn serializer_install_lock(
+ &self,
+ realm: &SerializerRealmKey,
+ ) -> (u64, Arc>) {
+ let mut runtime_state = self.runtime_state.lock().unwrap();
+ let generation = runtime_state.ensure_serializer_generation(realm);
+ let install_lock = runtime_state
+ .serializer_install_locks
+ .entry(realm.clone())
+ .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
+ .clone();
+ (generation, install_lock)
+ }
+
+ fn serializer_install_is_current(
+ &self,
+ realm: &SerializerRealmKey,
+ generation: u64,
+ install_lock: &Arc>,
+ ) -> bool {
+ let runtime_state = self.runtime_state.lock().unwrap();
+ runtime_state.serializer_generations.get(realm) == Some(&generation)
+ && runtime_state
+ .serializer_install_locks
+ .get(realm)
+ .is_some_and(|current| Arc::ptr_eq(current, install_lock))
+ }
+
+ fn remember_serializer(
+ &self,
+ realm: &SerializerRealmKey,
+ generation: u64,
+ install_lock: &Arc>,
+ object_id: String,
+ execution_context_id: Option,
+ ) -> bool {
+ let mut runtime_state = self.runtime_state.lock().unwrap();
+ let generation_is_current =
+ runtime_state.serializer_generations.get(realm) == Some(&generation);
+ let lock_is_current = runtime_state
+ .serializer_install_locks
+ .get(realm)
+ .is_some_and(|current| Arc::ptr_eq(current, install_lock));
+ if !generation_is_current || !lock_is_current {
+ return false;
+ }
+ let execution_context_id =
+ execution_context_id.or_else(|| runtime_state.execution_context_for_realm(realm));
+ runtime_state.serializers.insert(
+ realm.clone(),
+ SerializerCacheEntry {
+ object_id,
+ execution_context_id,
+ },
+ );
+ true
+ }
+
+ fn cleanup_serializer_install(
+ &self,
+ realm: &SerializerRealmKey,
+ generation: u64,
+ install_lock: &Arc>,
+ ) {
+ let mut runtime_state = self.runtime_state.lock().unwrap();
+ let exact_lock = runtime_state
+ .serializer_install_locks
+ .get(realm)
+ .is_some_and(|current| Arc::ptr_eq(current, install_lock));
+ if exact_lock {
+ runtime_state.serializer_install_locks.remove(realm);
+ }
+ if exact_lock
+ && runtime_state.serializer_generations.get(realm) == Some(&generation)
+ && !runtime_state.serializers.contains_key(realm)
+ {
+ runtime_state.serializer_generations.remove(realm);
+ }
+ }
+
+ fn release_uncommitted_serializer(&self, realm: &SerializerRealmKey, object_id: String) {
+ self.runtime_state
+ .lock()
+ .unwrap()
+ .enqueue_serializer_release(SerializerRelease {
+ session_id: realm.session_id.clone(),
+ object_id,
+ });
+ }
+
+ fn forget_serializer(&self, realm: &SerializerRealmKey) {
+ self.runtime_state
+ .lock()
+ .unwrap()
+ .evict_serializer_realms([realm.clone()], false);
+ }
+
+ fn frame_id_for_execution_context(
+ &self,
+ session_id: &str,
+ execution_context_id: &str,
+ ) -> Option {
+ self.runtime_state
+ .lock()
+ .unwrap()
+ .frame_id_for_execution_context(session_id, execution_context_id)
+ }
+
async fn connect(ws_endpoint: &str) -> RwResult> {
Self::connect_with_headers(ws_endpoint, &[]).await
}
@@ -17543,6 +18815,12 @@ impl CdpClient {
let traffic_log = Arc::new(Mutex::new(CdpTrafficLog::new()));
let traffic_log_writer = Arc::clone(&traffic_log);
let traffic_log_reader = Arc::clone(&traffic_log);
+ let (serializer_release_tx, serializer_release_rx) = mpsc::unbounded_channel();
+ let runtime_state = Arc::new(Mutex::new(CdpRuntimeState::new(Some(
+ serializer_release_tx,
+ ))));
+ let runtime_state_reader = Arc::clone(&runtime_state);
+ spawn_serializer_release_pump(serializer_release_rx, write_tx.clone());
let alive = Arc::new(AtomicBool::new(true));
let alive_writer = Arc::clone(&alive);
let alive_reader = Arc::clone(&alive);
@@ -17616,6 +18894,7 @@ impl CdpClient {
events_reader.clone(),
Arc::clone(&event_log_reader),
Arc::clone(&traffic_log_reader),
+ Arc::clone(&runtime_state_reader),
);
}
alive_reader.store(false, Ordering::SeqCst);
@@ -17630,6 +18909,7 @@ impl CdpClient {
events,
event_log,
traffic_log,
+ runtime_state,
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -17658,6 +18938,12 @@ impl CdpClient {
let traffic_log = Arc::new(Mutex::new(CdpTrafficLog::new()));
let traffic_log_writer = Arc::clone(&traffic_log);
let traffic_log_dispatcher = Arc::clone(&traffic_log);
+ let (serializer_release_tx, serializer_release_rx) = mpsc::unbounded_channel();
+ let runtime_state = Arc::new(Mutex::new(CdpRuntimeState::new(Some(
+ serializer_release_tx,
+ ))));
+ let runtime_state_dispatcher = Arc::clone(&runtime_state);
+ spawn_serializer_release_pump(serializer_release_rx, write_tx.clone());
let alive = Arc::new(AtomicBool::new(true));
let alive_writer = Arc::clone(&alive);
let alive_dispatcher = Arc::clone(&alive);
@@ -17745,6 +19031,7 @@ impl CdpClient {
events_dispatcher.clone(),
Arc::clone(&event_log_dispatcher),
Arc::clone(&traffic_log_dispatcher),
+ Arc::clone(&runtime_state_dispatcher),
);
}
alive_dispatcher.store(false, Ordering::SeqCst);
@@ -17763,6 +19050,7 @@ impl CdpClient {
events,
event_log,
traffic_log,
+ runtime_state,
next_id: AtomicU64::new(1),
sent_runtime_enable_count: AtomicU64::new(0),
sent_target_close_count: AtomicU64::new(0),
@@ -21162,8 +22450,8 @@ struct PyDownloadEventWaiter {
#[pyclass(name = "_FileChooserEventWaiter")]
struct PyFileChooserEventWaiter {
browser: Arc,
+ page: Arc,
receiver: Mutex