Improve page management and speed up execution (#87)
This commit is contained in:
@@ -44,16 +44,11 @@ class SyncSession:
|
||||
) -> PageInfo: # pragma: no cover
|
||||
"""Get a new page to use"""
|
||||
|
||||
# Close all finished pages to ensure clean state
|
||||
self.page_pool.close_all_finished_pages()
|
||||
|
||||
# If we're at max capacity after cleanup, wait for busy pages to finish
|
||||
if self.page_pool.pages_count >= self.max_pages:
|
||||
start_time = time()
|
||||
while time() - start_time < self._max_wait_for_page:
|
||||
# Wait for any pages to finish, then clean them up
|
||||
sleep(0.05)
|
||||
self.page_pool.close_all_finished_pages()
|
||||
if self.page_pool.pages_count < self.max_pages:
|
||||
break
|
||||
else:
|
||||
@@ -105,16 +100,11 @@ class AsyncSession(SyncSession):
|
||||
) -> PageInfo: # pragma: no cover
|
||||
"""Get a new page to use"""
|
||||
async with self._lock:
|
||||
# Close all finished pages to ensure clean state
|
||||
await self.page_pool.aclose_all_finished_pages()
|
||||
|
||||
# If we're at max capacity after cleanup, wait for busy pages to finish
|
||||
if self.page_pool.pages_count >= self.max_pages:
|
||||
start_time = time()
|
||||
while time() - start_time < self._max_wait_for_page:
|
||||
# Wait for any pages to finish, then clean them up
|
||||
await asyncio_sleep(0.05)
|
||||
await self.page_pool.aclose_all_finished_pages()
|
||||
if self.page_pool.pages_count < self.max_pages:
|
||||
break
|
||||
else:
|
||||
|
||||
@@ -381,8 +381,9 @@ class StealthySession(StealthySessionMixin, SyncSession):
|
||||
page_info.page, first_response, final_response, params.selector_config
|
||||
)
|
||||
|
||||
# Mark the page as finished for next use
|
||||
page_info.mark_finished()
|
||||
# Close the page, to free up resources
|
||||
page_info.page.close()
|
||||
self.page_pool.pages.remove(page_info)
|
||||
|
||||
return response
|
||||
|
||||
@@ -701,8 +702,9 @@ class AsyncStealthySession(StealthySessionMixin, AsyncSession):
|
||||
page_info.page, first_response, final_response, params.selector_config
|
||||
)
|
||||
|
||||
# Mark the page as finished for next use
|
||||
page_info.mark_finished()
|
||||
# Close the page, to free up resources
|
||||
await page_info.page.close()
|
||||
self.page_pool.pages.remove(page_info)
|
||||
|
||||
return response
|
||||
|
||||
|
||||
@@ -305,8 +305,9 @@ class DynamicSession(DynamicSessionMixin, SyncSession):
|
||||
page_info.page, first_response, final_response, params.selector_config
|
||||
)
|
||||
|
||||
# Mark the page as finished for next use
|
||||
page_info.mark_finished()
|
||||
# Close the page, to free up resources
|
||||
page_info.page.close()
|
||||
self.page_pool.pages.remove(page_info)
|
||||
|
||||
return response
|
||||
|
||||
@@ -554,9 +555,9 @@ class AsyncDynamicSession(DynamicSessionMixin, AsyncSession):
|
||||
page_info.page, first_response, final_response, params.selector_config
|
||||
)
|
||||
|
||||
# Mark the page as finished for next use
|
||||
page_info.mark_finished()
|
||||
|
||||
# Close the page, to free up resources
|
||||
await page_info.page.close()
|
||||
self.page_pool.pages.remove(page_info)
|
||||
return response
|
||||
|
||||
except Exception as e: # pragma: no cover
|
||||
|
||||
@@ -23,11 +23,6 @@ class PageInfo:
|
||||
self.state = "busy"
|
||||
self.url = url
|
||||
|
||||
def mark_finished(self):
|
||||
"""Mark the page as finished for new requests"""
|
||||
self.state = "finished"
|
||||
self.url = ""
|
||||
|
||||
def mark_error(self):
|
||||
"""Mark the page as having an error"""
|
||||
self.state = "error"
|
||||
@@ -83,33 +78,3 @@ class PagePool:
|
||||
"""Remove pages in error state"""
|
||||
with self._lock:
|
||||
self.pages = [p for p in self.pages if p.state != "error"]
|
||||
|
||||
def close_all_finished_pages(self):
|
||||
"""Close all pages in finished state and remove them from the pool"""
|
||||
with self._lock:
|
||||
pages_to_remove = []
|
||||
for page_info in self.pages:
|
||||
if page_info.state == "finished":
|
||||
try:
|
||||
page_info.page.close()
|
||||
except Exception:
|
||||
pass
|
||||
pages_to_remove.append(page_info)
|
||||
|
||||
for page_info in pages_to_remove:
|
||||
self.pages.remove(page_info)
|
||||
|
||||
async def aclose_all_finished_pages(self):
|
||||
"""Async version: Close all pages in finished state and remove them from the pool"""
|
||||
with self._lock:
|
||||
pages_to_remove = []
|
||||
for page_info in self.pages:
|
||||
if page_info.state == "finished":
|
||||
try:
|
||||
await page_info.page.close()
|
||||
except Exception:
|
||||
pass
|
||||
pages_to_remove.append(page_info)
|
||||
|
||||
for page_info in pages_to_remove:
|
||||
self.pages.remove(page_info)
|
||||
|
||||
@@ -54,16 +54,18 @@ class TestAsyncStealthySession:
|
||||
"""Test page pool creation and reuse"""
|
||||
async with AsyncStealthySession() as session:
|
||||
# The first request creates a page
|
||||
_ = await session.fetch(urls["basic"])
|
||||
assert session.page_pool.pages_count == 1
|
||||
response = await session.fetch(urls["basic"])
|
||||
assert response.status == 200
|
||||
assert session.page_pool.pages_count == 0
|
||||
|
||||
# The second request should reuse the page
|
||||
_ = await session.fetch(urls["html"])
|
||||
assert session.page_pool.pages_count == 1
|
||||
response = await session.fetch(urls["html"])
|
||||
assert response.status == 200
|
||||
assert session.page_pool.pages_count == 0
|
||||
|
||||
# Check pool stats
|
||||
stats = session.get_pool_stats()
|
||||
assert stats["total_pages"] == 1
|
||||
assert stats["total_pages"] == 0
|
||||
assert stats["max_pages"] == 1
|
||||
|
||||
async def test_stealthy_session_with_options(self, urls):
|
||||
|
||||
@@ -53,16 +53,18 @@ class TestAsyncDynamicSession:
|
||||
"""Test page pool creation and reuse"""
|
||||
async with AsyncDynamicSession() as session:
|
||||
# The first request creates a page
|
||||
_ = await session.fetch(urls["basic"])
|
||||
assert session.page_pool.pages_count == 1
|
||||
|
||||
response = await session.fetch(urls["basic"])
|
||||
assert response.status == 200
|
||||
assert session.page_pool.pages_count == 0
|
||||
|
||||
# The second request should reuse the page
|
||||
_ = await session.fetch(urls["html"])
|
||||
assert session.page_pool.pages_count == 1
|
||||
response = await session.fetch(urls["html"])
|
||||
assert response.status == 200
|
||||
assert session.page_pool.pages_count == 0
|
||||
|
||||
# Check pool stats
|
||||
stats = session.get_pool_stats()
|
||||
assert stats["total_pages"] == 1
|
||||
assert stats["total_pages"] == 0
|
||||
assert stats["max_pages"] == 1
|
||||
|
||||
async def test_dynamic_session_with_options(self, urls):
|
||||
|
||||
@@ -24,10 +24,6 @@ class TestPageInfo:
|
||||
assert page_info.state == "busy"
|
||||
assert page_info.url == "https://example.com"
|
||||
|
||||
page_info.mark_finished()
|
||||
assert page_info.state == "finished"
|
||||
assert page_info.url == ""
|
||||
|
||||
page_info.mark_error()
|
||||
assert page_info.state == "error"
|
||||
|
||||
@@ -88,21 +84,7 @@ class TestPagePool:
|
||||
with pytest.raises(RuntimeError):
|
||||
pool.add_page(Mock())
|
||||
|
||||
def test_get_ready_page(self):
|
||||
"""Test getting ready page"""
|
||||
pool = PagePool(max_pages=3)
|
||||
|
||||
# Add pages
|
||||
page1 = pool.add_page(Mock())
|
||||
page2 = pool.add_page(Mock())
|
||||
|
||||
# Mark them as finished
|
||||
page1.mark_finished()
|
||||
page2.mark_finished()
|
||||
|
||||
# test
|
||||
pool.close_all_finished_pages()
|
||||
assert pool.pages_count == 0
|
||||
|
||||
def test_cleanup_error_pages(self):
|
||||
"""Test cleaning up error pages"""
|
||||
|
||||
Reference in New Issue
Block a user