mirror of
https://github.com/msitarzewski/agency-agents.git
synced 2026-10-01 13:55:54 +03:00
Merge commit 'refs/c35/949' into integ/content-batch
This commit is contained in:
@@ -333,104 +333,138 @@ export class PageGeometryEngine {
|
||||
Maintains a warm Chromium browser instance with pooled, isolated `BrowserContext` objects, concurrency rate limiting, route blocking for external noise, and scheduled recycling to deliver sub-80ms compilations:
|
||||
|
||||
```python
|
||||
# cv_pdf_pool.py: High-Throughput Browser Context Pool
|
||||
# cv_pdf_pool.py: Drain active jobs before recycling the shared browser.
|
||||
import asyncio
|
||||
import logging
|
||||
from typing import Optional
|
||||
from playwright.async_api import async_playwright, Browser, BrowserContext, Playwright
|
||||
from playwright.async_api import async_playwright, Browser, Playwright
|
||||
|
||||
logger = logging.getLogger("pdf_pool")
|
||||
|
||||
class PlaywrightPDFPool:
|
||||
def __init__(self, max_concurrency: int = 4, max_jobs_before_recycle: int = 500):
|
||||
self.max_concurrency = max_concurrency
|
||||
if max_concurrency < 1 or max_jobs_before_recycle < 1:
|
||||
raise ValueError("Pool limits must be positive")
|
||||
self.max_jobs_before_recycle = max_jobs_before_recycle
|
||||
self.semaphore = asyncio.Semaphore(max_concurrency)
|
||||
self.job_counter = 0
|
||||
self.playwright: Optional[Playwright] = None
|
||||
self.browser: Optional[Browser] = None
|
||||
self._lock = asyncio.Lock()
|
||||
self._condition = asyncio.Condition()
|
||||
self._active_jobs = 0
|
||||
self._recycling = False
|
||||
self._closed = False
|
||||
self._shutdown_task = None
|
||||
|
||||
async def initialize(self):
|
||||
async with self._lock:
|
||||
if self.browser and self.browser.is_connected():
|
||||
return
|
||||
self.playwright = await async_playwright().start()
|
||||
async def _close_locked(self):
|
||||
browser, playwright = self.browser, self.playwright
|
||||
self.browser = self.playwright = None
|
||||
try:
|
||||
if browser:
|
||||
await browser.close()
|
||||
finally:
|
||||
if playwright:
|
||||
await playwright.stop()
|
||||
|
||||
async def _initialize_locked(self):
|
||||
if self.browser and self.browser.is_connected():
|
||||
return
|
||||
if self._active_jobs:
|
||||
raise RuntimeError("Disconnected browser still has active jobs")
|
||||
await self._close_locked()
|
||||
self.playwright = await async_playwright().start()
|
||||
try:
|
||||
self.browser = await self.playwright.chromium.launch(
|
||||
headless=True,
|
||||
args=[
|
||||
"--disable-background-networking",
|
||||
"--disable-gpu",
|
||||
"--disable-dev-shm-usage",
|
||||
"--no-sandbox",
|
||||
"--font-render-hinting=none"
|
||||
]
|
||||
args=["--disable-background-networking", "--disable-gpu",
|
||||
"--disable-dev-shm-usage", "--no-sandbox", "--font-render-hinting=none"]
|
||||
)
|
||||
self.job_counter = 0
|
||||
logger.info("Playwright PDF Pool initialized with warm Chromium instance.")
|
||||
except BaseException:
|
||||
await self._close_locked()
|
||||
raise
|
||||
self.job_counter = 0
|
||||
|
||||
async def render_pdf(
|
||||
self,
|
||||
html_content: str,
|
||||
width_mm: float = 210.0,
|
||||
height_mm: float = 297.0
|
||||
) -> bytes:
|
||||
await self.initialize()
|
||||
async def initialize(self):
|
||||
async with self._condition:
|
||||
await self._condition.wait_for(lambda: not self._recycling or self._closed)
|
||||
if self._closed:
|
||||
raise RuntimeError("PDF pool is shut down")
|
||||
await self._initialize_locked()
|
||||
|
||||
async def render_pdf(self, html_content: str, width_mm: float = 210.0,
|
||||
height_mm: float = 297.0) -> bytes:
|
||||
async with self.semaphore:
|
||||
self.job_counter += 1
|
||||
if self.job_counter >= self.max_jobs_before_recycle:
|
||||
logger.info("Recycling browser process after %d jobs.", self.job_counter)
|
||||
await self.recycle()
|
||||
|
||||
# Create isolated context for the request
|
||||
context: BrowserContext = await self.browser.new_context(
|
||||
viewport={"width": int(width_mm * 96 / 25.4), "height": int(height_mm * 96 / 25.4)},
|
||||
device_scale_factor=1.0
|
||||
)
|
||||
async with self._condition:
|
||||
# At the threshold, stop admitting new jobs until this
|
||||
# generation drains; sustained traffic cannot starve recycling.
|
||||
await self._condition.wait_for(lambda: self._closed or (
|
||||
not self._recycling and (self.job_counter < self.max_jobs_before_recycle
|
||||
or self._active_jobs == 0)
|
||||
))
|
||||
if self._closed:
|
||||
raise RuntimeError("PDF pool is shut down")
|
||||
# Drain at the configured threshold without interrupting active
|
||||
# renders or reacquiring a lock held by the same coroutine.
|
||||
if self._active_jobs == 0 and self.job_counter >= self.max_jobs_before_recycle:
|
||||
await self._close_locked()
|
||||
await self._initialize_locked()
|
||||
browser = self.browser
|
||||
self._active_jobs += 1
|
||||
self.job_counter += 1
|
||||
|
||||
context = None
|
||||
try:
|
||||
page = await context.new_page()
|
||||
|
||||
# Abort tracking and off-target external requests
|
||||
await page.route(
|
||||
"**/*",
|
||||
lambda route: route.abort() if route.request.resource_type in ["media", "websocket"] else route.continue_()
|
||||
context = await browser.new_context(
|
||||
viewport={"width": int(width_mm * 96 / 25.4), "height": int(height_mm * 96 / 25.4)},
|
||||
device_scale_factor=1.0
|
||||
)
|
||||
|
||||
# Load HTML with networkidle guarantee
|
||||
page = await context.new_page()
|
||||
await page.route("**/*", lambda route: route.abort()
|
||||
if route.request.resource_type in ["media", "websocket"]
|
||||
else route.continue_())
|
||||
await page.set_content(html_content, wait_until="networkidle")
|
||||
await page.evaluate("document.fonts.ready")
|
||||
|
||||
# Generate tagged, vector-clean PDF via CDP
|
||||
pdf_bytes = await page.pdf(
|
||||
width=f"{width_mm}mm",
|
||||
height=f"{height_mm}mm",
|
||||
print_background=True,
|
||||
prefer_css_page_size=True,
|
||||
tagged=True,
|
||||
return await page.pdf(
|
||||
width=f"{width_mm}mm", height=f"{height_mm}mm",
|
||||
print_background=True, prefer_css_page_size=True, tagged=True,
|
||||
margin={"top": "0mm", "right": "0mm", "bottom": "0mm", "left": "0mm"}
|
||||
)
|
||||
return pdf_bytes
|
||||
finally:
|
||||
await context.close()
|
||||
try:
|
||||
if context:
|
||||
await context.close()
|
||||
finally:
|
||||
async with self._condition:
|
||||
self._active_jobs -= 1
|
||||
self._condition.notify_all()
|
||||
|
||||
async def recycle(self):
|
||||
async with self._lock:
|
||||
if self.browser:
|
||||
await self.browser.close()
|
||||
if self.playwright:
|
||||
await self.playwright.stop()
|
||||
self.browser = None
|
||||
self.playwright = None
|
||||
await self.initialize()
|
||||
async with self._condition:
|
||||
await self._condition.wait_for(lambda: not self._recycling or self._closed)
|
||||
if self._closed:
|
||||
raise RuntimeError("PDF pool is shut down")
|
||||
self._recycling = True
|
||||
try:
|
||||
await self._condition.wait_for(lambda: self._active_jobs == 0)
|
||||
await self._close_locked()
|
||||
await self._initialize_locked()
|
||||
finally:
|
||||
self._recycling = False
|
||||
self._condition.notify_all()
|
||||
|
||||
async def _finish_shutdown(self):
|
||||
async with self._condition:
|
||||
self._condition.notify_all()
|
||||
await self._condition.wait_for(lambda: self._active_jobs == 0 and not self._recycling)
|
||||
await self._close_locked()
|
||||
|
||||
async def shutdown(self):
|
||||
async with self._lock:
|
||||
if self.browser:
|
||||
await self.browser.close()
|
||||
if self.playwright:
|
||||
await self.playwright.stop()
|
||||
# Retain cleanup: cancelling a caller must not abandon the browser
|
||||
# after in-flight renders finish.
|
||||
self._closed = True
|
||||
if self._shutdown_task is None:
|
||||
self._shutdown_task = asyncio.create_task(self._finish_shutdown())
|
||||
await asyncio.shield(self._shutdown_task)
|
||||
```
|
||||
|
||||
### 4. 1:1 Sheet Canvas Viewport Scaler Architecture (CSS & React)
|
||||
|
||||
Reference in New Issue
Block a user