mirror of
https://github.com/browser-use/browser-use.git
synced 2026-10-02 04:04:36 +08:00
refactor: implement error handling for asyncio tasks across various components
- Introduced a utility function `create_task_with_error_handling` to manage asyncio tasks with proper exception logging and suppression options. - Updated multiple components, including `SessionManager`, `BrowserSession`, and various watchdogs, to utilize the new error handling mechanism for background tasks. - Enhanced overall robustness of asynchronous operations by ensuring exceptions are logged and handled appropriately, improving maintainability and debugging capabilities.
This commit is contained in:
@@ -44,7 +44,7 @@ from browser_use.browser.profile import BrowserProfile, ProxySettings
|
||||
from browser_use.browser.views import BrowserStateSummary, TabInfo
|
||||
from browser_use.dom.views import DOMRect, EnhancedDOMTreeNode, TargetInfo
|
||||
from browser_use.observability import observe_debug
|
||||
from browser_use.utils import _log_pretty_url, is_new_tab_page
|
||||
from browser_use.utils import _log_pretty_url, create_task_with_error_handling, is_new_tab_page
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from browser_use.actor.page import Page
|
||||
@@ -1606,7 +1606,9 @@ class BrowserSession(BaseModel):
|
||||
self.logger.debug(f'Proxy auth respond failed: {type(e).__name__}: {e}')
|
||||
|
||||
# schedule
|
||||
asyncio.create_task(_respond())
|
||||
create_task_with_error_handling(
|
||||
_respond(), name='auth_respond', logger_instance=self.logger, suppress_exceptions=True
|
||||
)
|
||||
else:
|
||||
# Default behaviour for non-proxy challenges: let browser handle
|
||||
async def _default():
|
||||
@@ -1620,7 +1622,9 @@ class BrowserSession(BaseModel):
|
||||
self.logger.debug(f'Default auth respond failed: {type(e).__name__}: {e}')
|
||||
|
||||
if request_id:
|
||||
asyncio.create_task(_default())
|
||||
create_task_with_error_handling(
|
||||
_default(), name='auth_default', logger_instance=self.logger, suppress_exceptions=True
|
||||
)
|
||||
|
||||
def _on_request_paused(event: RequestPausedEvent, session_id: SessionID | None = None):
|
||||
# Continue all paused requests to avoid stalling the network
|
||||
@@ -1638,7 +1642,9 @@ class BrowserSession(BaseModel):
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
asyncio.create_task(_continue())
|
||||
create_task_with_error_handling(
|
||||
_continue(), name='request_continue', logger_instance=self.logger, suppress_exceptions=True
|
||||
)
|
||||
|
||||
# Register event handler on root client
|
||||
try:
|
||||
@@ -1670,7 +1676,9 @@ class BrowserSession(BaseModel):
|
||||
except Exception as e:
|
||||
self.logger.debug(f'Fetch.enable on attached session failed: {type(e).__name__}: {e}')
|
||||
|
||||
asyncio.create_task(_enable())
|
||||
create_task_with_error_handling(
|
||||
_enable(), name='fetch_enable_attached', logger_instance=self.logger, suppress_exceptions=True
|
||||
)
|
||||
|
||||
try:
|
||||
self._cdp_client_root.register.Target.attachedToTarget(_on_attached)
|
||||
|
||||
@@ -9,6 +9,8 @@ from typing import TYPE_CHECKING
|
||||
|
||||
from cdp_use.cdp.target import AttachedToTargetEvent, DetachedFromTargetEvent, SessionID, TargetID
|
||||
|
||||
from browser_use.utils import create_task_with_error_handling
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from browser_use.browser.session import BrowserSession, CDPSession, Target
|
||||
|
||||
@@ -72,14 +74,29 @@ class SessionManager:
|
||||
# - Create CDPSession
|
||||
# - Enable monitoring (for pages/tabs)
|
||||
# - Add to pool
|
||||
asyncio.create_task(self._handle_target_attached(event))
|
||||
create_task_with_error_handling(
|
||||
self._handle_target_attached(event),
|
||||
name='handle_target_attached',
|
||||
logger_instance=self.logger,
|
||||
suppress_exceptions=True,
|
||||
)
|
||||
|
||||
def on_detached(event: DetachedFromTargetEvent, session_id: SessionID | None = None):
|
||||
asyncio.create_task(self._handle_target_detached(event))
|
||||
create_task_with_error_handling(
|
||||
self._handle_target_detached(event),
|
||||
name='handle_target_detached',
|
||||
logger_instance=self.logger,
|
||||
suppress_exceptions=True,
|
||||
)
|
||||
|
||||
def on_target_info_changed(event, session_id: SessionID | None = None):
|
||||
# Update session info from targetInfoChanged events (no polling needed!)
|
||||
asyncio.create_task(self._handle_target_info_changed(event))
|
||||
create_task_with_error_handling(
|
||||
self._handle_target_info_changed(event),
|
||||
name='handle_target_info_changed',
|
||||
logger_instance=self.logger,
|
||||
suppress_exceptions=True,
|
||||
)
|
||||
|
||||
cdp_client.register.Target.attachedToTarget(on_attached)
|
||||
cdp_client.register.Target.detachedFromTarget(on_detached)
|
||||
@@ -629,7 +646,9 @@ class SessionManager:
|
||||
await asyncio.sleep(0.05)
|
||||
|
||||
# Start checking in background
|
||||
check_task = asyncio.create_task(check_all_ready())
|
||||
check_task = create_task_with_error_handling(
|
||||
check_all_ready(), name='check_all_targets_ready', logger_instance=self.logger
|
||||
)
|
||||
|
||||
try:
|
||||
# Wait for completion with timeout
|
||||
|
||||
@@ -18,6 +18,7 @@ from browser_use.browser.events import (
|
||||
TabCreatedEvent,
|
||||
)
|
||||
from browser_use.browser.watchdog_base import BaseWatchdog
|
||||
from browser_use.utils import create_task_with_error_handling
|
||||
|
||||
if TYPE_CHECKING:
|
||||
pass
|
||||
@@ -61,7 +62,9 @@ class CrashWatchdog(BaseWatchdog):
|
||||
"""Start monitoring when browser is connected."""
|
||||
# logger.debug('[CrashWatchdog] Browser connected event received, beginning monitoring')
|
||||
|
||||
asyncio.create_task(self._start_monitoring())
|
||||
create_task_with_error_handling(
|
||||
self._start_monitoring(), name='start_crash_monitoring', logger_instance=self.logger, suppress_exceptions=True
|
||||
)
|
||||
# logger.debug(f'[CrashWatchdog] Monitoring task started: {self._monitoring_task and not self._monitoring_task.done()}')
|
||||
|
||||
async def on_BrowserStoppedEvent(self, event: BrowserStoppedEvent) -> None:
|
||||
@@ -95,7 +98,12 @@ class CrashWatchdog(BaseWatchdog):
|
||||
# Register crash event handler
|
||||
def on_target_crashed(event: TargetCrashedEvent, session_id: SessionID | None = None):
|
||||
# Create and track the task
|
||||
task = asyncio.create_task(self._on_target_crash_cdp(target_id))
|
||||
task = create_task_with_error_handling(
|
||||
self._on_target_crash_cdp(target_id),
|
||||
name='handle_target_crash',
|
||||
logger_instance=self.logger,
|
||||
suppress_exceptions=True,
|
||||
)
|
||||
self._cdp_event_tasks.add(task)
|
||||
# Remove from set when done
|
||||
task.add_done_callback(lambda t: self._cdp_event_tasks.discard(t))
|
||||
@@ -194,7 +202,9 @@ class CrashWatchdog(BaseWatchdog):
|
||||
# logger.info('[CrashWatchdog] Monitoring already running')
|
||||
return
|
||||
|
||||
self._monitoring_task = asyncio.create_task(self._monitoring_loop())
|
||||
self._monitoring_task = create_task_with_error_handling(
|
||||
self._monitoring_loop(), name='crash_monitoring_loop', logger_instance=self.logger, suppress_exceptions=True
|
||||
)
|
||||
# logger.debug('[CrashWatchdog] Monitoring loop created and started')
|
||||
|
||||
async def _stop_monitoring(self) -> None:
|
||||
|
||||
@@ -17,7 +17,7 @@ from browser_use.dom.views import (
|
||||
SerializedDOMState,
|
||||
)
|
||||
from browser_use.observability import observe_debug
|
||||
from browser_use.utils import time_execution_async
|
||||
from browser_use.utils import create_task_with_error_handling, time_execution_async
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from browser_use.browser.views import BrowserStateSummary, NetworkRequest, PageInfo, PaginationButton
|
||||
@@ -369,12 +369,16 @@ class DOMWatchdog(BaseWatchdog):
|
||||
else None
|
||||
)
|
||||
|
||||
dom_task = asyncio.create_task(self._build_dom_tree_without_highlights(previous_state))
|
||||
dom_task = create_task_with_error_handling(
|
||||
self._build_dom_tree_without_highlights(previous_state), name='build_dom_tree', logger_instance=self.logger
|
||||
)
|
||||
|
||||
# Start clean screenshot task if requested (without JS highlights)
|
||||
if event.include_screenshot:
|
||||
self.logger.debug('🔍 DOMWatchdog.on_BrowserStateRequestEvent: 📸 Starting clean screenshot task...')
|
||||
screenshot_task = asyncio.create_task(self._capture_clean_screenshot())
|
||||
screenshot_task = create_task_with_error_handling(
|
||||
self._capture_clean_screenshot(), name='capture_screenshot', logger_instance=self.logger
|
||||
)
|
||||
|
||||
# Wait for both tasks to complete
|
||||
content = None
|
||||
|
||||
@@ -25,6 +25,7 @@ from browser_use.browser.events import (
|
||||
TabCreatedEvent,
|
||||
)
|
||||
from browser_use.browser.watchdog_base import BaseWatchdog
|
||||
from browser_use.utils import create_task_with_error_handling
|
||||
|
||||
if TYPE_CHECKING:
|
||||
pass
|
||||
@@ -182,7 +183,12 @@ class DownloadsWatchdog(BaseWatchdog):
|
||||
except (AssertionError, KeyError):
|
||||
pass
|
||||
# Create and track the task
|
||||
task = asyncio.create_task(self._handle_cdp_download(event, target_id, session_id))
|
||||
task = create_task_with_error_handling(
|
||||
self._handle_cdp_download(event, target_id, session_id),
|
||||
name='handle_cdp_download',
|
||||
logger_instance=self.logger,
|
||||
suppress_exceptions=True,
|
||||
)
|
||||
self._cdp_event_tasks.add(task)
|
||||
# Remove from set when done
|
||||
task.add_done_callback(lambda t: self._cdp_event_tasks.discard(t))
|
||||
@@ -431,7 +437,12 @@ class DownloadsWatchdog(BaseWatchdog):
|
||||
self.logger.error(f'[DownloadsWatchdog] Error downloading in background: {type(e).__name__}: {e}')
|
||||
|
||||
# Create background task
|
||||
task = asyncio.create_task(download_in_background())
|
||||
task = create_task_with_error_handling(
|
||||
download_in_background(),
|
||||
name='download_in_background',
|
||||
logger_instance=self.logger,
|
||||
suppress_exceptions=True,
|
||||
)
|
||||
self._cdp_event_tasks.add(task)
|
||||
task.add_done_callback(lambda t: self._cdp_event_tasks.discard(t))
|
||||
|
||||
|
||||
@@ -12,6 +12,7 @@ from browser_use.browser.events import BrowserConnectedEvent, BrowserStopEvent
|
||||
from browser_use.browser.profile import ViewportSize
|
||||
from browser_use.browser.video_recorder import VideoRecorderService
|
||||
from browser_use.browser.watchdog_base import BaseWatchdog
|
||||
from browser_use.utils import create_task_with_error_handling
|
||||
|
||||
|
||||
class RecordingWatchdog(BaseWatchdog):
|
||||
@@ -97,10 +98,16 @@ class RecordingWatchdog(BaseWatchdog):
|
||||
"""
|
||||
Synchronous handler for incoming screencast frames.
|
||||
"""
|
||||
|
||||
if not self._recorder:
|
||||
return
|
||||
self._recorder.add_frame(event['data'])
|
||||
asyncio.create_task(self._ack_screencast_frame(event, session_id))
|
||||
create_task_with_error_handling(
|
||||
self._ack_screencast_frame(event, session_id),
|
||||
name='ack_screencast_frame',
|
||||
logger_instance=self.logger,
|
||||
suppress_exceptions=True,
|
||||
)
|
||||
|
||||
async def _ack_screencast_frame(self, event: ScreencastFrameEvent, session_id: str | None) -> None:
|
||||
"""
|
||||
|
||||
@@ -19,6 +19,7 @@ from browser_use.browser.events import (
|
||||
StorageStateSavedEvent,
|
||||
)
|
||||
from browser_use.browser.watchdog_base import BaseWatchdog
|
||||
from browser_use.utils import create_task_with_error_handling
|
||||
|
||||
|
||||
class StorageStateWatchdog(BaseWatchdog):
|
||||
@@ -91,7 +92,9 @@ class StorageStateWatchdog(BaseWatchdog):
|
||||
|
||||
assert self.browser_session.cdp_client is not None
|
||||
|
||||
self._monitoring_task = asyncio.create_task(self._monitor_storage_changes())
|
||||
self._monitoring_task = create_task_with_error_handling(
|
||||
self._monitor_storage_changes(), name='monitor_storage_changes', logger_instance=self.logger, suppress_exceptions=True
|
||||
)
|
||||
# self.logger'[StorageStateWatchdog] Started storage monitoring task')
|
||||
|
||||
async def _stop_monitoring(self) -> None:
|
||||
|
||||
+8
-2
@@ -1854,8 +1854,14 @@ async def run_auth_command():
|
||||
|
||||
# Run authentication and progress updates concurrently
|
||||
auth_start_time = asyncio.get_event_loop().time()
|
||||
auth_task = asyncio.create_task(sync_service.authenticate(show_instructions=True))
|
||||
progress_task = asyncio.create_task(show_auth_progress())
|
||||
from browser_use.utils import create_task_with_error_handling
|
||||
|
||||
auth_task = create_task_with_error_handling(
|
||||
sync_service.authenticate(show_instructions=True), name='sync_authenticate'
|
||||
)
|
||||
progress_task = create_task_with_error_handling(
|
||||
show_auth_progress(), name='show_auth_progress', suppress_exceptions=True
|
||||
)
|
||||
|
||||
# Wait for authentication to complete, with timeout
|
||||
success = await asyncio.wait_for(auth_task, timeout=120.0) # 2 minutes for initial testing
|
||||
|
||||
@@ -23,6 +23,7 @@ from browser_use.dom.views import (
|
||||
TargetAllTrees,
|
||||
)
|
||||
from browser_use.observability import observe_debug
|
||||
from browser_use.utils import create_task_with_error_handling
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from browser_use.browser.session import BrowserSession
|
||||
@@ -317,10 +318,10 @@ class DomService:
|
||||
|
||||
# Create initial tasks
|
||||
tasks = {
|
||||
'snapshot': asyncio.create_task(create_snapshot_request()),
|
||||
'dom_tree': asyncio.create_task(create_dom_tree_request()),
|
||||
'ax_tree': asyncio.create_task(self._get_ax_tree_for_all_frames(target_id)),
|
||||
'device_pixel_ratio': asyncio.create_task(self._get_viewport_ratio(target_id)),
|
||||
'snapshot': create_task_with_error_handling(create_snapshot_request(), name='get_snapshot'),
|
||||
'dom_tree': create_task_with_error_handling(create_dom_tree_request(), name='get_dom_tree'),
|
||||
'ax_tree': create_task_with_error_handling(self._get_ax_tree_for_all_frames(target_id), name='get_ax_tree'),
|
||||
'device_pixel_ratio': create_task_with_error_handling(self._get_viewport_ratio(target_id), name='get_viewport_ratio'),
|
||||
}
|
||||
|
||||
# Wait for all tasks with timeout
|
||||
@@ -333,10 +334,14 @@ class DomService:
|
||||
|
||||
# Retry mapping for pending tasks
|
||||
retry_map = {
|
||||
tasks['snapshot']: lambda: asyncio.create_task(create_snapshot_request()),
|
||||
tasks['dom_tree']: lambda: asyncio.create_task(create_dom_tree_request()),
|
||||
tasks['ax_tree']: lambda: asyncio.create_task(self._get_ax_tree_for_all_frames(target_id)),
|
||||
tasks['device_pixel_ratio']: lambda: asyncio.create_task(self._get_viewport_ratio(target_id)),
|
||||
tasks['snapshot']: lambda: create_task_with_error_handling(create_snapshot_request(), name='get_snapshot_retry'),
|
||||
tasks['dom_tree']: lambda: create_task_with_error_handling(create_dom_tree_request(), name='get_dom_tree_retry'),
|
||||
tasks['ax_tree']: lambda: create_task_with_error_handling(
|
||||
self._get_ax_tree_for_all_frames(target_id), name='get_ax_tree_retry'
|
||||
),
|
||||
tasks['device_pixel_ratio']: lambda: create_task_with_error_handling(
|
||||
self._get_viewport_ratio(target_id), name='get_viewport_ratio_retry'
|
||||
),
|
||||
}
|
||||
|
||||
# Create new tasks only for the ones that didn't complete
|
||||
|
||||
@@ -33,7 +33,7 @@ from browser_use.agent.views import ActionResult
|
||||
from browser_use.telemetry import MCPClientTelemetryEvent, ProductTelemetry
|
||||
from browser_use.tools.registry.service import Registry
|
||||
from browser_use.tools.service import Tools
|
||||
from browser_use.utils import get_browser_use_version
|
||||
from browser_use.utils import create_task_with_error_handling, get_browser_use_version
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -93,7 +93,9 @@ class MCPClient:
|
||||
server_params = StdioServerParameters(command=self.command, args=self.args, env=self.env)
|
||||
|
||||
# Start stdio client in background task
|
||||
self._stdio_task = asyncio.create_task(self._run_stdio_client(server_params))
|
||||
self._stdio_task = create_task_with_error_handling(
|
||||
self._run_stdio_client(server_params), name='mcp_stdio_client', suppress_exceptions=True
|
||||
)
|
||||
|
||||
# Wait for connection to be established
|
||||
retries = 0
|
||||
|
||||
@@ -150,7 +150,7 @@ except ImportError:
|
||||
sys.exit(1)
|
||||
|
||||
from browser_use.telemetry import MCPServerTelemetryEvent, ProductTelemetry
|
||||
from browser_use.utils import get_browser_use_version
|
||||
from browser_use.utils import create_task_with_error_handling, get_browser_use_version
|
||||
|
||||
|
||||
def get_parent_process_cmdline() -> str | None:
|
||||
@@ -1059,7 +1059,7 @@ class BrowserUseServer:
|
||||
logger.error(f'Error in cleanup task: {e}')
|
||||
await asyncio.sleep(120)
|
||||
|
||||
self._cleanup_task = asyncio.create_task(cleanup_loop())
|
||||
self._cleanup_task = create_task_with_error_handling(cleanup_loop(), name='mcp_cleanup_loop', suppress_exceptions=True)
|
||||
|
||||
async def run(self):
|
||||
"""Run the MCP server."""
|
||||
|
||||
@@ -5,7 +5,6 @@ Fetches pricing data from LiteLLM repository and caches it for 1 day.
|
||||
Automatically tracks token usage when LLMs are registered and invoked.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
from datetime import datetime, timedelta
|
||||
@@ -29,6 +28,7 @@ from browser_use.tokens.views import (
|
||||
TokenUsageEntry,
|
||||
UsageSummary,
|
||||
)
|
||||
from browser_use.utils import create_task_with_error_handling
|
||||
|
||||
load_dotenv()
|
||||
|
||||
@@ -347,7 +347,9 @@ class TokenCost:
|
||||
|
||||
logger.debug(f'Token cost service: {usage}')
|
||||
|
||||
asyncio.create_task(token_cost_service._log_usage(llm.model, usage))
|
||||
create_task_with_error_handling(
|
||||
token_cost_service._log_usage(llm.model, usage), name='log_token_usage', suppress_exceptions=True
|
||||
)
|
||||
|
||||
# else:
|
||||
# await token_cost_service._log_non_usage_llm(llm)
|
||||
|
||||
@@ -51,7 +51,7 @@ from browser_use.tools.views import (
|
||||
SwitchTabAction,
|
||||
UploadFileAction,
|
||||
)
|
||||
from browser_use.utils import time_execution_sync
|
||||
from browser_use.utils import create_task_with_error_handling, time_execution_sync
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -257,7 +257,9 @@ class Tools(Generic[Context]):
|
||||
element_desc = get_click_description(node)
|
||||
|
||||
# Highlight the element being clicked (truly non-blocking)
|
||||
asyncio.create_task(browser_session.highlight_interaction_element(node))
|
||||
create_task_with_error_handling(
|
||||
browser_session.highlight_interaction_element(node), name='highlight_click_element', suppress_exceptions=True
|
||||
)
|
||||
|
||||
event = browser_session.event_bus.dispatch(ClickElementEvent(node=node))
|
||||
await event
|
||||
@@ -312,7 +314,9 @@ class Tools(Generic[Context]):
|
||||
return ActionResult(extracted_content=msg)
|
||||
|
||||
# Highlight the element being typed into (truly non-blocking)
|
||||
asyncio.create_task(browser_session.highlight_interaction_element(node))
|
||||
create_task_with_error_handling(
|
||||
browser_session.highlight_interaction_element(node), name='highlight_type_element', suppress_exceptions=True
|
||||
)
|
||||
|
||||
# Dispatch type text event with node
|
||||
try:
|
||||
@@ -459,7 +463,11 @@ class Tools(Generic[Context]):
|
||||
|
||||
# Highlight the file input element if found (truly non-blocking)
|
||||
if file_input_node:
|
||||
asyncio.create_task(browser_session.highlight_interaction_element(file_input_node))
|
||||
create_task_with_error_handling(
|
||||
browser_session.highlight_interaction_element(file_input_node),
|
||||
name='highlight_file_input',
|
||||
suppress_exceptions=True,
|
||||
)
|
||||
|
||||
# If not found near the selected element, fallback to finding the closest file input to current scroll position
|
||||
if file_input_node is None:
|
||||
@@ -494,8 +502,13 @@ class Tools(Generic[Context]):
|
||||
if closest_file_input:
|
||||
file_input_node = closest_file_input
|
||||
logger.info(f'Found file input closest to scroll position (distance: {min_distance}px)')
|
||||
|
||||
# Highlight the fallback file input element (truly non-blocking)
|
||||
asyncio.create_task(browser_session.highlight_interaction_element(file_input_node))
|
||||
create_task_with_error_handling(
|
||||
browser_session.highlight_interaction_element(file_input_node),
|
||||
name='highlight_file_input_fallback',
|
||||
suppress_exceptions=True,
|
||||
)
|
||||
else:
|
||||
msg = 'No file upload element found on the page'
|
||||
logger.error(msg)
|
||||
@@ -1616,7 +1629,11 @@ class CodeAgentTools(Tools[Context]):
|
||||
|
||||
# Highlight the file input element if found (truly non-blocking)
|
||||
if file_input_node:
|
||||
asyncio.create_task(browser_session.highlight_interaction_element(file_input_node))
|
||||
create_task_with_error_handling(
|
||||
browser_session.highlight_interaction_element(file_input_node),
|
||||
name='highlight_file_input',
|
||||
suppress_exceptions=True,
|
||||
)
|
||||
|
||||
# If not found near the selected element, fallback to finding the closest file input to current scroll position
|
||||
if file_input_node is None:
|
||||
@@ -1651,8 +1668,13 @@ class CodeAgentTools(Tools[Context]):
|
||||
if closest_file_input:
|
||||
file_input_node = closest_file_input
|
||||
logger.info(f'Found file input closest to scroll position (distance: {min_distance}px)')
|
||||
|
||||
# Highlight the fallback file input element (truly non-blocking)
|
||||
asyncio.create_task(browser_session.highlight_interaction_element(file_input_node))
|
||||
create_task_with_error_handling(
|
||||
browser_session.highlight_interaction_element(file_input_node),
|
||||
name='highlight_file_input_fallback',
|
||||
suppress_exceptions=True,
|
||||
)
|
||||
else:
|
||||
msg = 'No file upload element found on the page'
|
||||
logger.error(msg)
|
||||
|
||||
@@ -668,3 +668,53 @@ def _log_pretty_url(s: str, max_len: int | None = 22) -> str:
|
||||
if max_len is not None and len(s) > max_len:
|
||||
return s[:max_len] + '…'
|
||||
return s
|
||||
|
||||
|
||||
def create_task_with_error_handling(
|
||||
coro: Coroutine[Any, Any, T],
|
||||
*,
|
||||
name: str | None = None,
|
||||
logger_instance: logging.Logger | None = None,
|
||||
suppress_exceptions: bool = False,
|
||||
) -> asyncio.Task[T]:
|
||||
"""
|
||||
Create an asyncio task with proper exception handling to prevent "Task exception was never retrieved" warnings.
|
||||
|
||||
Args:
|
||||
coro: The coroutine to wrap in a task
|
||||
name: Optional name for the task (useful for debugging)
|
||||
logger_instance: Optional logger instance to use. If None, uses module logger.
|
||||
suppress_exceptions: If True, exceptions are only logged and not re-raised. Default False.
|
||||
|
||||
Returns:
|
||||
asyncio.Task: The created task with exception handling callback
|
||||
|
||||
Example:
|
||||
# Instead of: asyncio.create_task(some_async_function())
|
||||
# Use: create_task_with_error_handling(some_async_function(), name="my_task")
|
||||
"""
|
||||
task = asyncio.create_task(coro, name=name)
|
||||
log = logger_instance or logger
|
||||
|
||||
def _handle_task_exception(t: asyncio.Task[T]) -> None:
|
||||
"""Callback to handle task exceptions"""
|
||||
try:
|
||||
# This will raise if the task had an exception
|
||||
exc = t.exception()
|
||||
if exc is not None:
|
||||
task_name = t.get_name() if hasattr(t, 'get_name') else 'unnamed'
|
||||
log.error(f'Exception in background task [{task_name}]: {type(exc).__name__}: {exc}', exc_info=exc)
|
||||
|
||||
# Re-raise if not suppressed (useful for critical tasks)
|
||||
if not suppress_exceptions:
|
||||
raise exc
|
||||
except asyncio.CancelledError:
|
||||
# Task was cancelled, this is normal behavior
|
||||
pass
|
||||
except Exception as e:
|
||||
# Catch any other exception during exception handling
|
||||
task_name = t.get_name() if hasattr(t, 'get_name') else 'unnamed'
|
||||
log.error(f'Error handling exception in task [{task_name}]: {type(e).__name__}: {e}')
|
||||
|
||||
task.add_done_callback(_handle_task_exception)
|
||||
return task
|
||||
|
||||
@@ -7,6 +7,8 @@ import sys
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
|
||||
from browser_use.utils import create_task_with_error_handling
|
||||
|
||||
|
||||
def setup_environment(debug: bool):
|
||||
if not debug:
|
||||
@@ -85,7 +87,9 @@ Return ONLY the key brand info, not page structure details.""",
|
||||
screenshot_path = self.output_dir / f'landing_page_{timestamp}.png'
|
||||
await agent_instance.browser_session.take_screenshot(path=str(screenshot_path), full_page=False)
|
||||
|
||||
screenshot_task = asyncio.create_task(screenshot_callback(agent))
|
||||
screenshot_task = create_task_with_error_handling(
|
||||
screenshot_callback(agent), name='screenshot_callback', suppress_exceptions=True
|
||||
)
|
||||
history = await agent.run()
|
||||
try:
|
||||
await screenshot_task
|
||||
@@ -373,7 +377,7 @@ async def create_multiple_ads(url: str, debug: bool = False, mode: str = 'instag
|
||||
|
||||
tasks = []
|
||||
for i in range(count):
|
||||
task = asyncio.create_task(generate_single_ad(page_data, mode, i + 1))
|
||||
task = create_task_with_error_handling(generate_single_ad(page_data, mode, i + 1), name=f'generate_ad_{i + 1}')
|
||||
tasks.append(task)
|
||||
|
||||
results = await asyncio.gather(*tasks, return_exceptions=True)
|
||||
|
||||
+20
-33
@@ -59,26 +59,11 @@ dependencies = [
|
||||
# textual: used for terminal UI
|
||||
|
||||
[project.optional-dependencies]
|
||||
cli = [
|
||||
"textual>=3.2.0",
|
||||
]
|
||||
code = [
|
||||
"matplotlib>=3.9.0",
|
||||
"numpy>=2.3.2",
|
||||
"pandas>=2.2.0",
|
||||
"tabulate>=0.9.0",
|
||||
|
||||
]
|
||||
aws = [
|
||||
"boto3>=1.38.45"
|
||||
]
|
||||
oci = [
|
||||
"oci>=2.126.4",
|
||||
]
|
||||
video = [
|
||||
"imageio[ffmpeg]>=2.37.0",
|
||||
"numpy>=2.3.2",
|
||||
]
|
||||
cli = ["textual>=3.2.0"]
|
||||
code = ["matplotlib>=3.9.0", "numpy>=2.3.2", "pandas>=2.2.0", "tabulate>=0.9.0"]
|
||||
aws = ["boto3>=1.38.45"]
|
||||
oci = ["oci>=2.126.4"]
|
||||
video = ["imageio[ffmpeg]>=2.37.0", "numpy>=2.3.2"]
|
||||
examples = [
|
||||
"agentmail==0.0.59",
|
||||
# botocore: only needed for Bedrock Claude boto3 examples/models/bedrock_claude.py
|
||||
@@ -92,14 +77,10 @@ eval = [
|
||||
"lmnr[all]==0.7.17",
|
||||
"anyio>=4.9.0",
|
||||
"psutil>=7.0.0",
|
||||
"datamodel-code-generator>=0.26.0"
|
||||
]
|
||||
cli-oci = [
|
||||
"browser-use[cli,oci]",
|
||||
]
|
||||
all = [
|
||||
"browser-use[cli,examples,aws,oci]",
|
||||
"datamodel-code-generator>=0.26.0",
|
||||
]
|
||||
cli-oci = ["browser-use[cli,oci]"]
|
||||
all = ["browser-use[cli,examples,aws,oci]"]
|
||||
|
||||
# will prefer to use local source code checked out in ../../browser-use (if present) instead of pypi browser-use package
|
||||
# [tool.uv.sources]
|
||||
@@ -128,7 +109,15 @@ fix = true
|
||||
|
||||
[tool.ruff.lint]
|
||||
select = ["ASYNC", "E", "F", "FAST", "I", "PLE"]
|
||||
ignore = ["ASYNC109", "E101", "E402", "E501", "F841", "E731", "W291"] # TODO: determine if adding timeouts to all the unbounded async functions is needed / worth-it so we can un-ignore ASYNC109
|
||||
ignore = [
|
||||
"ASYNC109",
|
||||
"E101",
|
||||
"E402",
|
||||
"E501",
|
||||
"F841",
|
||||
"E731",
|
||||
"W291",
|
||||
] # TODO: determine if adding timeouts to all the unbounded async functions is needed / worth-it so we can un-ignore ASYNC109
|
||||
unfixable = ["E101", "E402", "E501", "F841", "E731"]
|
||||
|
||||
[tool.ruff.format]
|
||||
@@ -161,7 +150,7 @@ exclude = [
|
||||
"browser_use/llm/tests/test_single_step.py",
|
||||
"product_extraction.py",
|
||||
"discover/",
|
||||
"list/"
|
||||
"list/",
|
||||
]
|
||||
venvPath = "."
|
||||
venv = ".venv"
|
||||
@@ -195,9 +184,7 @@ markers = [
|
||||
"unit: marks tests as unit tests",
|
||||
"asyncio: mark tests as async tests",
|
||||
]
|
||||
testpaths = [
|
||||
"tests"
|
||||
]
|
||||
testpaths = ["tests"]
|
||||
python_files = ["test_*.py", "*_test.py"]
|
||||
addopts = "-svx --strict-markers --tb=short --dist=loadscope"
|
||||
log_cli = true
|
||||
@@ -241,5 +228,5 @@ dev-dependencies = [
|
||||
"lmnr[all]==0.7.17",
|
||||
# "pytest-playwright-asyncio>=0.7.0", # not actually needed I think
|
||||
"pytest-timeout>=2.4.0",
|
||||
"pydantic_settings>=2.10.1"
|
||||
"pydantic_settings>=2.10.1",
|
||||
]
|
||||
|
||||
Reference in New Issue
Block a user