mirror of
https://github.com/VectifyAI/PageIndex.git
synced 2026-07-18 21:21:05 +02:00
fix(index): bound LLM concurrency at the leaf, not per gather call
The previous approach (bounded_gather building a fresh semaphore per call) did NOT compose: the indexing call graph nests gathers (tree_parser -> process_large_node_recursively -> recurse, plus each node's check_title gather), and each level got its own independent cap. Peak in-flight LLM calls grew ~N^depth, so a deep/wide document still exhausted file descriptors (Errno 24) — the exact failure the cap was meant to prevent — while the flat summary phase was over-serialized. Naively sharing one semaphore across gather levels would instead deadlock (a parent holds a slot while awaiting children that need slots). Move the throttle to the single chokepoint every LLM call funnels through, llm_acompletion: one shared semaphore per event loop, acquired only around the litellm.acompletion network call. This gives a true global cap that composes across any nesting and can't deadlock (a parent awaiting children holds no slot). bounded_gather is gone; the call sites revert to plain asyncio.gather. Also: - Reject bool in max_concurrency validation (bool is an int subclass, so set_max_concurrency(True) / IndexConfig(max_concurrency=True) previously became Semaphore(1) and silently serialized). Shared _validate_max_concurrency + a pydantic field_validator. - Guard check_title_appearance_in_start_concurrent against an out-of-range or 0 physical_index (LLM can emit one): it was dereferenced during task construction, outside the gather's return_exceptions protection, aborting the whole build; 0 silently wrapped to the last page. Now marked 'no'. - Propagate contextvars into agent.py's worker-thread run (mirrors pipeline._run_async) so ContextVar settings stay consistent. - Tests rewritten to cover the nested case the old flat tests missed, the leaf-level throttle in llm_acompletion, bool rejection, and the out-of-range physical_index guard. Claude-Session: https://claude.ai/code/session_01Kx5DgKbhK1N8autqXH8SmS
This commit is contained in:
parent
e5392836a4
commit
2d46d68052
7 changed files with 203 additions and 130 deletions
|
|
@ -159,6 +159,10 @@ class AgentRunner:
|
|||
result = Runner.run_sync(agent, question)
|
||||
else:
|
||||
import concurrent.futures
|
||||
import contextvars
|
||||
# Copy the current context into the worker thread so ContextVar-based
|
||||
# settings propagate (mirrors pipeline._run_async).
|
||||
ctx = contextvars.copy_context()
|
||||
with concurrent.futures.ThreadPoolExecutor(max_workers=1) as pool:
|
||||
result = pool.submit(asyncio.run, Runner.run(agent, question)).result()
|
||||
result = pool.submit(ctx.run, asyncio.run, Runner.run(agent, question)).result()
|
||||
return result.final_output
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue