From 589f0b6708bd264cf93ecbcb858ce68fb4460091 Mon Sep 17 00:00:00 2001 From: nightcityblade Date: Thu, 13 Aug 2026 00:06:06 +0800 Subject: [PATCH] fix: close dispatcher stream in arun_many --- crawl4ai/async_webcrawler.py | 15 +++++---- tests/async/test_arun_many_stream_cleanup.py | 34 ++++++++++++++++++++ 2 files changed, 43 insertions(+), 6 deletions(-) create mode 100644 tests/async/test_arun_many_stream_cleanup.py diff --git a/crawl4ai/async_webcrawler.py b/crawl4ai/async_webcrawler.py index 8216d19bc..917d702df 100644 --- a/crawl4ai/async_webcrawler.py +++ b/crawl4ai/async_webcrawler.py @@ -9,7 +9,7 @@ import asyncio # from contextlib import nullcontext, asynccontextmanager -from contextlib import asynccontextmanager +from contextlib import aclosing, asynccontextmanager from .models import ( CrawlResult, MarkdownGenerationResult, @@ -1108,10 +1108,13 @@ async def maybe_release_session(): if stream: async def result_transformer(): try: - async for task_result in dispatcher.run_urls_stream( - crawler=self, urls=urls, config=config - ): - yield transform_result(task_result) + async with aclosing( + dispatcher.run_urls_stream( + crawler=self, urls=urls, config=config + ) + ) as task_results: + async for task_result in task_results: + yield transform_result(task_result) finally: # Auto-release session after streaming completes await maybe_release_session() @@ -1246,4 +1249,4 @@ async def amap_domain( config or DomainMapperConfig(**kwargs) if kwargs else DomainMapperConfig() ) - return await self._domain_mapper.scan(domain, mapper_config) \ No newline at end of file + return await self._domain_mapper.scan(domain, mapper_config) diff --git a/tests/async/test_arun_many_stream_cleanup.py b/tests/async/test_arun_many_stream_cleanup.py new file mode 100644 index 000000000..79584e3fc --- /dev/null +++ b/tests/async/test_arun_many_stream_cleanup.py @@ -0,0 +1,34 @@ +from types import SimpleNamespace + +import pytest + +from crawl4ai import AsyncWebCrawler, CrawlerRunConfig + + +class ClosingDispatcher: + closed = False + + async def run_urls_stream(self, **kwargs): + try: + yield SimpleNamespace( + result=SimpleNamespace(), + task_id="task", + memory_usage=0, + peak_memory=0, + start_time=0, + end_time=0, + error_message="", + ) + finally: + self.closed = True + + +@pytest.mark.asyncio +async def test_arun_many_closes_dispatcher_stream(): + dispatcher = ClosingDispatcher() + stream = await AsyncWebCrawler().arun_many( + ["url"], config=CrawlerRunConfig(stream=True), dispatcher=dispatcher + ) + await stream.__anext__() + await stream.aclose() + assert dispatcher.closed