diff --git a/backend/danswer/connectors/web/connector.py b/backend/danswer/connectors/web/connector.py index 50b34dac4b8..fb3eae87ab1 100644 --- a/backend/danswer/connectors/web/connector.py +++ b/backend/danswer/connectors/web/connector.py @@ -51,6 +51,43 @@ ) +def is_cloudflare_challenge(status: int | None, headers: dict[str, str]) -> bool: + """Did Cloudflare serve its bot check ("Just a moment...") in place of the + page? It marks those responses with `cf-mitigated: challenge` (403, or 503 + for the older JS challenge). The challenge page's JS also rewrites page.url + to a one-time `?__cf_chl_rt_tk=` URL, which must not be treated as a + redirect.""" + return ( + status in (403, 503) and headers.get("cf-mitigated", "").lower() == "challenge" + ) + + +def fetch_html_without_browser(url: str) -> tuple[str, str] | None: + """Fallback for pages whose Cloudflare bot check blocks headless Chromium. + + Cloudflare challenges clients that claim to be a browser but fail its + fingerprinting (headless Chromium, or anything sending DEFAULT_USER_AGENT), + while many sites let a plainly declared non-browser client through — e.g. + w3.org serves the real page to python-requests but a 403 challenge to our + Chrome user agent. So this deliberately sends requests' default user agent. + + Returns (final_url, html), or None if the fetch fails or isn't HTML.""" + try: + response = requests.get(url, timeout=WEB_CONNECTOR_PAGE_TIMEOUT_MS / 1000) + except requests.RequestException as e: + logger.warning(f"Non-browser fetch of {url} failed: {e}") + return None + + content_type = response.headers.get("content-type", "") + if not response.ok or "html" not in content_type: + logger.info( + f"Non-browser fetch of {url} returned HTTP {response.status_code} " + f"({content_type or 'no content-type'})" + ) + return None + return response.url, response.text + + def _is_browser_dead(exc: Exception) -> bool: """Heuristic: did the exception kill the browser/context (vs. just this page)? Only then is a full Playwright restart warranted; otherwise we retry @@ -162,31 +199,36 @@ def get_internal_links( def start_playwright() -> Tuple[Playwright, BrowserContext]: playwright = sync_playwright().start() - # Optionally launch an installed browser (e.g. system Chrome) instead of the - # bundled Chromium. Prod leaves this unset and uses the pinned Chromium in - # the Linux image; it exists only so local dev can fall back to a working - # browser when the pinned build is incompatible with the host OS (the - # chromium-1097 build for playwright 1.41.2 SIGSEGVs on recent macOS). - browser_channel = os.environ.get("WEB_CONNECTOR_BROWSER_CHANNEL") or None - browser = playwright.chromium.launch(headless=True, channel=browser_channel) - - context = browser.new_context(user_agent=DEFAULT_USER_AGENT) - - if ( - WEB_CONNECTOR_OAUTH_CLIENT_ID - and WEB_CONNECTOR_OAUTH_CLIENT_SECRET - and WEB_CONNECTOR_OAUTH_TOKEN_URL - ): - client = BackendApplicationClient(client_id=WEB_CONNECTOR_OAUTH_CLIENT_ID) - oauth = OAuth2Session(client=client) - token = oauth.fetch_token( - token_url=WEB_CONNECTOR_OAUTH_TOKEN_URL, - client_id=WEB_CONNECTOR_OAUTH_CLIENT_ID, - client_secret=WEB_CONNECTOR_OAUTH_CLIENT_SECRET, - ) - context.set_extra_http_headers( - {"Authorization": "Bearer {}".format(token["access_token"])} - ) + try: + # Optionally launch an installed browser (e.g. system Chrome) instead of the + # bundled Chromium. Prod leaves this unset and uses the pinned Chromium in + # the Linux image; it exists only so local dev can fall back to a working + # browser when the pinned build is incompatible with the host OS (the + # chromium-1097 build for playwright 1.41.2 SIGSEGVs on recent macOS). + browser_channel = os.environ.get("WEB_CONNECTOR_BROWSER_CHANNEL") or None + browser = playwright.chromium.launch(headless=True, channel=browser_channel) + + context = browser.new_context(user_agent=DEFAULT_USER_AGENT) + + if ( + WEB_CONNECTOR_OAUTH_CLIENT_ID + and WEB_CONNECTOR_OAUTH_CLIENT_SECRET + and WEB_CONNECTOR_OAUTH_TOKEN_URL + ): + client = BackendApplicationClient(client_id=WEB_CONNECTOR_OAUTH_CLIENT_ID) + oauth = OAuth2Session(client=client) + token = oauth.fetch_token( + token_url=WEB_CONNECTOR_OAUTH_TOKEN_URL, + client_id=WEB_CONNECTOR_OAUTH_CLIENT_ID, + client_secret=WEB_CONNECTOR_OAUTH_CLIENT_SECRET, + ) + context.set_extra_http_headers( + {"Authorization": "Bearer {}".format(token["access_token"])} + ) + except Exception: + # Don't leak a started instance on this thread (see load_from_state). + playwright.stop() + raise return playwright, context @@ -509,161 +551,203 @@ def load_from_state(self, is_polling: bool = False) -> GenerateDocumentsOutput: check_internet_connection(base_url) playwright, context = start_playwright() - restart_playwright = False - pages_visited = 0 - while to_visit: - if WEB_CONNECTOR_MAX_PAGES and pages_visited >= WEB_CONNECTOR_MAX_PAGES: - logger.info( - f"Reached WEB_CONNECTOR_MAX_PAGES ({WEB_CONNECTOR_MAX_PAGES}); " - f"stopping crawl with {len(to_visit)} URL(s) still queued." - ) - break + # Always stop the current instance, however the crawl ends (0 docs -> + # RuntimeError, an unexpected exception, or the caller abandoning this + # generator). An instance that isn't stopped leaves its asyncio loop + # marked running on this thread, so every later start_playwright() on + # the thread fails with "using Playwright Sync API inside the asyncio + # loop". Dask workers run a single thread each, so one leak broke every + # later web connector on that pod until it restarted. + try: + restart_playwright = False + pages_visited = 0 + while to_visit: + if WEB_CONNECTOR_MAX_PAGES and pages_visited >= WEB_CONNECTOR_MAX_PAGES: + logger.info( + f"Reached WEB_CONNECTOR_MAX_PAGES ({WEB_CONNECTOR_MAX_PAGES}); " + f"stopping crawl with {len(to_visit)} URL(s) still queued." + ) + break - current_url = to_visit.pop() - if current_url in visited_links: - continue - visited_links.add(current_url) + current_url = to_visit.pop() + if current_url in visited_links: + continue + visited_links.add(current_url) - try: - protected_url_check(current_url) - except Exception as e: - last_error = f"Invalid URL {current_url} due to {e}" - logger.warning(last_error) - continue - - logger.info(f"Visiting {current_url}") - pages_visited += 1 - - # Reinit the browser if a previous batch/crash flagged it. Done for - # every page (as before) so the browser is always live when we reach - # the batch-yield / final-stop below. - if restart_playwright: - playwright, context = start_playwright() - restart_playwright = False - - # --- PDF: fetch directly (no browser). timeout so a hung download - # can't stall the whole attempt. Matches the original: PDFs don't - # trigger the per-batch flush. --- - if current_url.split(".")[-1] == "pdf": try: - response = requests.get(current_url, timeout=60) - response.raise_for_status() - page_text = pdf_to_text(file=io.BytesIO(response.content)) - doc_batch.append( - Document( - id=current_url, - sections=[Section(link=current_url, text=page_text)], - source=DocumentSource.WEB, - semantic_identifier=current_url.split(".")[-1], - metadata={}, - ) - ) + protected_url_check(current_url) except Exception as e: - last_error = f"Failed to fetch PDF '{current_url}': {e}" - logger.error(last_error) - continue - - # --- HTML via Playwright, with retries. A single page error retries - # with a FRESH PAGE on the same browser (exponential backoff); only - # a browser-level crash restarts Playwright. One bad page no longer - # tears down the browser or fails the attempt. --- - page_doc: Document | None = None - for attempt in range(WEB_CONNECTOR_MAX_RETRIES): - if attempt > 0: - time.sleep(min(2**attempt + random.uniform(0, 1), 10)) - try: - page = context.new_page() + last_error = f"Invalid URL {current_url} due to {e}" + logger.warning(last_error) + continue + + logger.info(f"Visiting {current_url}") + pages_visited += 1 + + # Reinit the browser if a previous batch/crash flagged it. Done for + # every page (as before) so the browser is always live when we reach + # the batch-yield / final-stop below. + if restart_playwright: + playwright, context = start_playwright() + restart_playwright = False + + # --- PDF: fetch directly (no browser). timeout so a hung download + # can't stall the whole attempt. Matches the original: PDFs don't + # trigger the per-batch flush. --- + if current_url.split(".")[-1] == "pdf": try: - page_response = page.goto( - current_url, - timeout=WEB_CONNECTOR_PAGE_TIMEOUT_MS, - # 'domcontentloaded' (DOM parsed) instead of the - # default 'load' (waits for every image/font/etc.) — - # far faster and enough for text extraction. - wait_until="domcontentloaded", + response = requests.get(current_url, timeout=60) + response.raise_for_status() + page_text = pdf_to_text(file=io.BytesIO(response.content)) + doc_batch.append( + Document( + id=current_url, + sections=[Section(link=current_url, text=page_text)], + source=DocumentSource.WEB, + semantic_identifier=current_url.split(".")[-1], + metadata={}, + ) ) - final_page = page.url - if final_page != current_url: - logger.info(f"Redirected to {final_page}") - protected_url_check(final_page) - current_url = final_page - if current_url in visited_links: - logger.info("Redirected page already indexed") - break - visited_links.add(current_url) - - content = page.content() - soup = BeautifulSoup(content, "html.parser") - - if self.recursive and not is_polling: - for link in get_internal_links( - base_url, + except Exception as e: + last_error = f"Failed to fetch PDF '{current_url}': {e}" + logger.error(last_error) + continue + + # --- HTML via Playwright, with retries. A single page error retries + # with a FRESH PAGE on the same browser (exponential backoff); only + # a browser-level crash restarts Playwright. One bad page no longer + # tears down the browser or fails the attempt. --- + page_doc: Document | None = None + for attempt in range(WEB_CONNECTOR_MAX_RETRIES): + if attempt > 0: + time.sleep(min(2**attempt + random.uniform(0, 1), 10)) + try: + page = context.new_page() + try: + page_response = page.goto( current_url, - soup, - allowed_prefixes=self.recursive_prefixes, + timeout=WEB_CONNECTOR_PAGE_TIMEOUT_MS, + # 'domcontentloaded' (DOM parsed) instead of the + # default 'load' (waits for every image/font/etc.) — + # far faster and enough for text extraction. + wait_until="domcontentloaded", + ) + page_status = ( + page_response.status if page_response else None + ) + fallback_html: str | None = None + if page_response and is_cloudflare_challenge( + page_response.status, page_response.headers ): - if link not in visited_links: - to_visit.append(link) - - if page_response and str(page_response.status)[0] in ( - "4", - "5", - ): - last_error = ( - f"Skipped indexing {current_url} due to HTTP " - f"{page_response.status} response" + # Retrying the browser just gets the challenge again; + # try once without it instead. + logger.info( + f"Cloudflare bot challenge on {current_url}; " + "retrying without the browser" + ) + fallback = fetch_html_without_browser(current_url) + if fallback is None: + last_error = ( + f"Skipped indexing {current_url}: blocked by a " + f"Cloudflare bot challenge (HTTP {page_status}), " + "and a non-browser fetch also failed" + ) + logger.info(last_error) + break + final_page, fallback_html = fallback + page_status = 200 + else: + final_page = page.url + if final_page != current_url: + logger.info(f"Redirected to {final_page}") + protected_url_check(final_page) + current_url = final_page + if current_url in visited_links: + logger.info("Redirected page already indexed") + break + visited_links.add(current_url) + + content = ( + fallback_html + if fallback_html is not None + else page.content() + ) + soup = BeautifulSoup(content, "html.parser") + + if self.recursive and not is_polling: + for link in get_internal_links( + base_url, + current_url, + soup, + allowed_prefixes=self.recursive_prefixes, + ): + if link not in visited_links: + to_visit.append(link) + + if page_status and str(page_status)[0] in ("4", "5"): + last_error = ( + f"Skipped indexing {current_url} due to HTTP " + f"{page_status} response" + ) + logger.info(last_error) + break # a real 4xx/5xx — don't retry + + parsed_html = web_html_cleanup(soup, self.mintlify_cleanup) + page_doc = Document( + id=current_url, + sections=[ + Section( + link=current_url, text=parsed_html.cleaned_text + ) + ], + source=DocumentSource.WEB, + semantic_identifier=parsed_html.title or current_url, + metadata={}, ) - logger.info(last_error) - break # a real 4xx/5xx — don't retry - - parsed_html = web_html_cleanup(soup, self.mintlify_cleanup) - page_doc = Document( - id=current_url, - sections=[ - Section(link=current_url, text=parsed_html.cleaned_text) - ], - source=DocumentSource.WEB, - semantic_identifier=parsed_html.title or current_url, - metadata={}, + break # success + finally: + page.close() + except Exception as e: + last_error = ( + f"Failed to fetch '{current_url}' " + f"(attempt {attempt + 1}/{WEB_CONNECTOR_MAX_RETRIES}): {e}" ) - break # success - finally: - page.close() - except Exception as e: - last_error = ( - f"Failed to fetch '{current_url}' " - f"(attempt {attempt + 1}/{WEB_CONNECTOR_MAX_RETRIES}): {e}" - ) - logger.warning(last_error) - if _is_browser_dead(e): - # Browser/context crashed — restart it so the next - # attempt (and subsequent pages) have a live browser. - try: - playwright.stop() - except Exception: - pass - playwright, context = start_playwright() - # else: transient page error — retry with a fresh page. - - if page_doc is not None: - doc_batch.append(page_doc) - - if len(doc_batch) >= self.batch_size: + logger.warning(last_error) + if _is_browser_dead(e): + # Browser/context crashed — restart it so the next + # attempt (and subsequent pages) have a live browser. + try: + playwright.stop() + except Exception: + pass + playwright, context = start_playwright() + # else: transient page error — retry with a fresh page. + + if page_doc is not None: + doc_batch.append(page_doc) + + if len(doc_batch) >= self.batch_size: + playwright.stop() + restart_playwright = True + at_least_one_doc = True + yield doc_batch + doc_batch = [] + + if doc_batch: playwright.stop() - restart_playwright = True at_least_one_doc = True yield doc_batch - doc_batch = [] - - if doc_batch: - playwright.stop() - at_least_one_doc = True - yield doc_batch - if not at_least_one_doc: - if last_error: - raise RuntimeError(last_error) - raise RuntimeError("No valid pages found.") + if not at_least_one_doc: + if last_error: + raise RuntimeError(last_error) + raise RuntimeError("No valid pages found.") + finally: + # No-op if already stopped (batch boundary / final batch). + try: + playwright.stop() + except Exception as e: + logger.warning(f"Failed to stop Playwright: {e}") def poll_source( self, start: SecondsSinceUnixEpoch, end: SecondsSinceUnixEpoch diff --git a/backend/tests/unit/danswer/connectors/web/test_cloudflare_fallback.py b/backend/tests/unit/danswer/connectors/web/test_cloudflare_fallback.py new file mode 100644 index 00000000000..a059dbe41df --- /dev/null +++ b/backend/tests/unit/danswer/connectors/web/test_cloudflare_fallback.py @@ -0,0 +1,105 @@ +"""Tests for the web connector's Cloudflare bot-challenge fallback. + +Cloudflare serves a 403 "Just a moment..." page (marked `cf-mitigated: +challenge`) to headless Chromium on some sites (e.g. w3.org), and its JS rewrites +page.url to a one-time `?__cf_chl_rt_tk=` URL. The connector detects that and +re-fetches the page without the browser. These guard the detection and the +fallback's accept/reject rules. +""" +from typing import Any + +import pytest +import requests + +from danswer.connectors.web import connector as web_connector +from danswer.connectors.web.connector import fetch_html_without_browser +from danswer.connectors.web.connector import is_cloudflare_challenge + + +class _FakeResponse: + def __init__( + self, + status_code: int = 200, + content_type: str = "text/html; charset=utf-8", + text: str = "Real page", + url: str = "https://www.w3.org/TR/WCAG22/", + ) -> None: + self.status_code = status_code + self.ok = status_code < 400 + self.headers = {"content-type": content_type} if content_type else {} + self.text = text + self.url = url + + +def _patch_get(monkeypatch: pytest.MonkeyPatch, result: Any) -> list[dict]: + calls: list[dict] = [] + + def fake_get(url: str, **kwargs: Any) -> Any: + calls.append({"url": url, **kwargs}) + if isinstance(result, Exception): + raise result + return result + + monkeypatch.setattr(web_connector.requests, "get", fake_get) + return calls + + +@pytest.mark.parametrize("status", [403, 503]) +def test_detects_cloudflare_challenge(status: int) -> None: + assert is_cloudflare_challenge(status, {"cf-mitigated": "challenge"}) + + +def test_challenge_header_value_is_case_insensitive() -> None: + assert is_cloudflare_challenge(403, {"cf-mitigated": "Challenge"}) + + +@pytest.mark.parametrize( + "status,headers", + [ + (403, {}), # plain 403 — a real "forbidden", not a bot check + (403, {"server": "cloudflare"}), # Cloudflare-fronted but not mitigated + (200, {"cf-mitigated": "challenge"}), + (404, {"cf-mitigated": "challenge"}), + (None, {"cf-mitigated": "challenge"}), + ], +) +def test_ignores_non_challenge_responses( + status: int | None, headers: dict[str, str] +) -> None: + assert not is_cloudflare_challenge(status, headers) + + +def test_fallback_returns_final_url_and_html(monkeypatch: pytest.MonkeyPatch) -> None: + calls = _patch_get( + monkeypatch, + _FakeResponse(url="https://www.w3.org/TR/WCAG22/Overview.html"), + ) + result = fetch_html_without_browser("https://www.w3.org/TR/WCAG22/") + assert result == ( + "https://www.w3.org/TR/WCAG22/Overview.html", + "Real page", + ) + # Must not send the browser user agent — that's what Cloudflare challenges. + assert "headers" not in calls[0] + assert calls[0]["timeout"] > 0 + + +@pytest.mark.parametrize( + "response", + [ + _FakeResponse(status_code=403), # still challenged / blocked + _FakeResponse(status_code=500), + _FakeResponse(content_type="application/pdf"), + _FakeResponse(content_type=""), + ], +) +def test_fallback_rejects_unusable_responses( + monkeypatch: pytest.MonkeyPatch, response: _FakeResponse +) -> None: + _patch_get(monkeypatch, response) + assert fetch_html_without_browser("https://www.w3.org/TR/WCAG22/") is None + + +def test_fallback_swallows_request_errors(monkeypatch: pytest.MonkeyPatch) -> None: + _patch_get(monkeypatch, requests.ConnectionError("boom")) + assert fetch_html_without_browser("https://www.w3.org/TR/WCAG22/") is None diff --git a/backend/tests/unit/danswer/connectors/web/test_playwright_cleanup.py b/backend/tests/unit/danswer/connectors/web/test_playwright_cleanup.py new file mode 100644 index 00000000000..f2156c47324 --- /dev/null +++ b/backend/tests/unit/danswer/connectors/web/test_playwright_cleanup.py @@ -0,0 +1,136 @@ +"""Tests that the web connector always stops Playwright. + +A sync_playwright instance that isn't stopped leaves its asyncio loop marked +running on the thread, so the next start_playwright() on that thread fails with +"using Playwright Sync API inside the asyncio loop". Dask workers run one thread +each, so a single leak broke every later web connector on the pod. These drive +load_from_state with a fake browser and check every started instance ends up +stopped, however the crawl ends. +""" +from typing import Any + +import pytest + +from danswer.connectors.web import connector as web_connector +from danswer.connectors.web.connector import start_playwright +from danswer.connectors.web.connector import WebConnector + + +class _FakeResponse: + def __init__(self, status: int) -> None: + self.status = status + self.headers: dict[str, str] = {} + + +class _FakePage: + def __init__(self, pages: dict[str, tuple[int, str]]) -> None: + self._pages = pages + self.url = "" + + def goto(self, url: str, **kwargs: Any) -> _FakeResponse: + self.url = url + return _FakeResponse(self._pages[url][0]) + + def content(self) -> str: + return self._pages[self.url][1] + + def close(self) -> None: + pass + + +class _FakePlaywright: + def __init__(self) -> None: + self.stopped = False + + def stop(self) -> None: + self.stopped = True + + +class _FakeContext: + def __init__(self, pages: dict[str, tuple[int, str]]) -> None: + self._pages = pages + + def new_page(self) -> _FakePage: + return _FakePage(self._pages) + + +def _fake_browser( + monkeypatch: pytest.MonkeyPatch, pages: dict[str, tuple[int, str]] +) -> list[_FakePlaywright]: + started: list[_FakePlaywright] = [] + + def fake_start_playwright() -> tuple[_FakePlaywright, _FakeContext]: + pw = _FakePlaywright() + started.append(pw) + return pw, _FakeContext(pages) + + monkeypatch.setattr(web_connector, "start_playwright", fake_start_playwright) + monkeypatch.setattr(web_connector, "check_internet_connection", lambda url: None) + return started + + +def _html(title: str, body: str = "") -> str: + return f"{title}{body}" + + +def test_stops_playwright_when_crawl_finds_no_documents( + monkeypatch: pytest.MonkeyPatch, +) -> None: + # The prod case: the only page is skipped (HTTP 4xx), so load_from_state + # raises RuntimeError — previously without stopping Playwright. + started = _fake_browser( + monkeypatch, {"https://site.test/": (403, _html("Just a moment..."))} + ) + connector = WebConnector(base_url="https://site.test/", web_connector_type="single") + + with pytest.raises(RuntimeError, match="HTTP 403"): + list(connector.load_from_state()) + + assert len(started) == 1 + assert started[0].stopped + + +def test_stops_instance_restarted_after_a_batch_when_later_pages_yield_nothing( + monkeypatch: pytest.MonkeyPatch, +) -> None: + # Page A fills a batch (instance 1 stopped, restart flagged); page B starts + # instance 2 but 404s, so the crawl ends with an empty batch and no error — + # previously instance 2 leaked silently. + started = _fake_browser( + monkeypatch, + { + "https://site.test/": (200, _html("A", 'b')), + "https://site.test/b": (404, _html("Not found")), + }, + ) + connector = WebConnector( + base_url="https://site.test/", web_connector_type="recursive", batch_size=1 + ) + + batches = list(connector.load_from_state()) + + assert [d.id for batch in batches for d in batch] == ["https://site.test/"] + assert len(started) == 2 + assert all(pw.stopped for pw in started) + + +def test_start_playwright_stops_instance_if_browser_launch_fails( + monkeypatch: pytest.MonkeyPatch, +) -> None: + pw = _FakePlaywright() + + class _FailingChromium: + def launch(self, **kwargs: Any) -> None: + raise RuntimeError("launch failed") + + pw.chromium = _FailingChromium() # type: ignore[attr-defined] + + class _FakeSyncPlaywright: + def start(self) -> _FakePlaywright: + return pw + + monkeypatch.setattr(web_connector, "sync_playwright", _FakeSyncPlaywright) + + with pytest.raises(RuntimeError, match="launch failed"): + start_playwright() + assert pw.stopped