Spaces:
Running
Run the blocking DuckDuckGo call in a worker thread
Browse files`DDGS.text()` is a synchronous method doing blocking HTTP I/O, and it was
being called directly from an `async def`. That froze the entire event loop
for the duration of every DuckDuckGo request: no other request progressed,
no in-flight Playwright scrape progressed, nothing else on the loop got a
turn.
DuckDuckGo is the first backend tried for every query in `/serp/search`,
which is the flagship MCP tool, and both the endpoint docstrings and
MCP_INSTRUCTIONS tell clients to batch their query variations into a single
call - so the advertised concurrency was not just absent on the primary
path, it was inverted into N sequential stalls.
Measured before the fix, with a 0.25s blocking stand-in: four gathered
queries took 1.00s (fully serialized) and the event loop got 0 turns.
After: the same four take about as long as one, and the loop keeps ticking
throughout.
tests/test_serp_ddg_concurrency.py asserts both properties directly - the
gathered queries don't serialize, and a heartbeat task keeps getting turns
while a query is in flight - so a future change that puts the blocking call
back on the loop fails the suite rather than silently costing throughput.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01MSNSYnceqvdVz7Csis4K9e
- serp.py +23 -6
- tests/test_serp_ddg_concurrency.py +72 -0
|
@@ -8,6 +8,7 @@ from urllib.parse import quote_plus
|
|
| 8 |
import logging
|
| 9 |
import re
|
| 10 |
from lxml import etree
|
|
|
|
| 11 |
from asyncio import Semaphore
|
| 12 |
|
| 13 |
# Concurrency limit for Playwright browser contexts.
|
|
@@ -277,14 +278,30 @@ async def query_bing_search(browser: Browser, q: str, n_results: int = 10):
|
|
| 277 |
return await _extract_bing_results(page, n_results)
|
| 278 |
|
| 279 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 280 |
async def query_ddg_search(q: str, n_results: int = 10):
|
| 281 |
-
"""Queries duckduckgo search for the specified query
|
| 282 |
-
|
| 283 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 284 |
|
| 285 |
-
|
| 286 |
-
|
| 287 |
-
|
|
|
|
| 288 |
|
| 289 |
if not results:
|
| 290 |
raise DuckDuckGoBlockedException()
|
|
|
|
| 8 |
import logging
|
| 9 |
import re
|
| 10 |
from lxml import etree
|
| 11 |
+
import asyncio
|
| 12 |
from asyncio import Semaphore
|
| 13 |
|
| 14 |
# Concurrency limit for Playwright browser contexts.
|
|
|
|
| 278 |
return await _extract_bing_results(page, n_results)
|
| 279 |
|
| 280 |
|
| 281 |
+
def _ddg_text_blocking(q: str, n_results: int) -> list[dict]:
|
| 282 |
+
"""The synchronous half of a DuckDuckGo query.
|
| 283 |
+
|
| 284 |
+
`DDGS.text()` performs blocking HTTP I/O. It is deliberately isolated
|
| 285 |
+
here so `query_ddg_search` can hand it to a worker thread rather than
|
| 286 |
+
calling it on the event loop.
|
| 287 |
+
"""
|
| 288 |
+
return list(DDGS().text(q, max_results=n_results))
|
| 289 |
+
|
| 290 |
+
|
| 291 |
async def query_ddg_search(q: str, n_results: int = 10):
|
| 292 |
+
"""Queries duckduckgo search for the specified query.
|
| 293 |
+
|
| 294 |
+
The underlying library call is synchronous, so it runs in a worker
|
| 295 |
+
thread: calling it directly would freeze the whole event loop for the
|
| 296 |
+
duration of the request, and DuckDuckGo is the first backend tried for
|
| 297 |
+
every query in `/serp/search`.
|
| 298 |
+
"""
|
| 299 |
+
raw_results = await asyncio.to_thread(_ddg_text_blocking, q, n_results)
|
| 300 |
|
| 301 |
+
results = [
|
| 302 |
+
{"title": r["title"], "body": r["body"], "href": r["href"]}
|
| 303 |
+
for r in raw_results
|
| 304 |
+
]
|
| 305 |
|
| 306 |
if not results:
|
| 307 |
raise DuckDuckGoBlockedException()
|
|
@@ -0,0 +1,72 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
"""The DuckDuckGo backend must not block the event loop.
|
| 2 |
+
|
| 3 |
+
`DDGS.text()` is a synchronous, blocking method. Calling it directly from
|
| 4 |
+
an `async def` freezes the whole process for the duration of the HTTP
|
| 5 |
+
request: no other request progresses, no in-flight Playwright scrape
|
| 6 |
+
progresses, and nothing else on the loop gets a turn.
|
| 7 |
+
|
| 8 |
+
That matters more here than it looks. DuckDuckGo is the *first* backend
|
| 9 |
+
tried for every query in `/serp/search`, which is the flagship MCP tool,
|
| 10 |
+
and both the endpoint docstrings and MCP_INSTRUCTIONS actively encourage
|
| 11 |
+
clients to batch many queries into one call - which under a blocking
|
| 12 |
+
implementation serializes into N sequential stalls instead of N concurrent
|
| 13 |
+
requests.
|
| 14 |
+
|
| 15 |
+
These tests use a blocking stand-in with a known duration rather than the
|
| 16 |
+
real library, so they assert the concurrency property itself and stay fast
|
| 17 |
+
and offline.
|
| 18 |
+
"""
|
| 19 |
+
|
| 20 |
+
import asyncio
|
| 21 |
+
import time
|
| 22 |
+
|
| 23 |
+
import serp
|
| 24 |
+
from serp import query_ddg_search
|
| 25 |
+
|
| 26 |
+
BLOCK_SECONDS = 0.25
|
| 27 |
+
|
| 28 |
+
|
| 29 |
+
class _BlockingDDGS:
|
| 30 |
+
"""Stands in for the real DDGS: a synchronous call that takes real time."""
|
| 31 |
+
|
| 32 |
+
def text(self, q, max_results=10):
|
| 33 |
+
time.sleep(BLOCK_SECONDS)
|
| 34 |
+
return [{"title": q, "body": "b", "href": "h"}]
|
| 35 |
+
|
| 36 |
+
|
| 37 |
+
async def test_concurrent_queries_do_not_serialize(monkeypatch):
|
| 38 |
+
"""Four gathered queries should take about as long as one, not four."""
|
| 39 |
+
monkeypatch.setattr(serp, "DDGS", lambda: _BlockingDDGS())
|
| 40 |
+
|
| 41 |
+
started = time.perf_counter()
|
| 42 |
+
await asyncio.gather(*[query_ddg_search(f"q{i}", 5) for i in range(4)])
|
| 43 |
+
elapsed = time.perf_counter() - started
|
| 44 |
+
|
| 45 |
+
assert elapsed < BLOCK_SECONDS * 2, (
|
| 46 |
+
f"4 concurrent queries took {elapsed:.2f}s; "
|
| 47 |
+
f"serialized would be ~{BLOCK_SECONDS * 4:.2f}s, concurrent ~{BLOCK_SECONDS:.2f}s")
|
| 48 |
+
|
| 49 |
+
|
| 50 |
+
async def test_the_event_loop_stays_responsive_during_a_query(monkeypatch):
|
| 51 |
+
"""Anything else scheduled on the loop must keep getting turns while a
|
| 52 |
+
DuckDuckGo query is in flight - this is what a blocked loop costs every
|
| 53 |
+
other request being served at the same time.
|
| 54 |
+
"""
|
| 55 |
+
monkeypatch.setattr(serp, "DDGS", lambda: _BlockingDDGS())
|
| 56 |
+
|
| 57 |
+
ticks = 0
|
| 58 |
+
stop = asyncio.Event()
|
| 59 |
+
|
| 60 |
+
async def heartbeat():
|
| 61 |
+
nonlocal ticks
|
| 62 |
+
while not stop.is_set():
|
| 63 |
+
await asyncio.sleep(0.01)
|
| 64 |
+
ticks += 1
|
| 65 |
+
|
| 66 |
+
beat = asyncio.create_task(heartbeat())
|
| 67 |
+
await query_ddg_search("widget", 5)
|
| 68 |
+
stop.set()
|
| 69 |
+
await beat
|
| 70 |
+
|
| 71 |
+
# A responsive loop ticks ~25 times in 0.25s; a blocked one manages ~1.
|
| 72 |
+
assert ticks > 10, f"event loop only got {ticks} turns during a {BLOCK_SECONDS}s query"
|