Skip to content

Commit cd23188

Browse files
author
Sebastian Braun
committed
fix(indexer): close PageIndex's local SQLite connection and forward LLM credentials
index_long_document() now closes PageIndex's local SQLite connection(s) in a finally block covering both the success and the failure path (via a new best-effort _close_pageindex_client() helper, reaching into the pinned pageindex version's private LocalBackend/SQLiteStorage since it exposes no public close()/context-manager API). Without this, a subsequent mutation rollback could not unlink/rename pageindex.db on Windows while this process still held it open (WinError 32), leaving a dirty mutation journal and aborting the rest of a multi-file 'openkb add' batch (#249). Also resolves the KB's LlmCredentialBundle (the same LLM_API_KEY/base_url compiler.py's own _llm_call/_llm_call_async use) and forwards it into PageIndex's own internal LLM calls via IndexConfig(llm_params=...) -- PageIndex's pinned version scopes llm_params per-call (pageindex.config.llm_params_scope, context-isolated, safe under concurrent multi-KB use) but had no wiring from OpenKB's side, so a custom LLM_API_KEY/gateway base_url never reached PageIndex's TOC/tree/summary generation calls, which instead fell back to LiteLLM's default provider-key/env-var lookup (#219). Guarded by the same IndexConfig.model_fields check already used for max_concurrency, so it degrades gracefully against an older pinned pageindex.
1 parent ff54396 commit cd23188

2 files changed

Lines changed: 273 additions & 86 deletions

File tree

‎openkb/indexer.py‎

Lines changed: 142 additions & 85 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@
1111

1212
from pageindex import IndexConfig, PageIndexClient
1313

14-
from openkb.config import resolve_concurrency, resolve_effective_config
14+
from openkb.config import resolve_concurrency, resolve_credential_bundle, resolve_effective_config
1515
from openkb.tree_renderer import render_summary_md
1616

1717
logger = logging.getLogger(__name__)
@@ -153,7 +153,7 @@ def _write_long_doc_artifacts(
153153
return summary_path
154154

155155

156-
def _build_index_config(config: dict[str, Any]) -> IndexConfig:
156+
def _build_index_config(config: dict[str, Any], bundle=None) -> IndexConfig:
157157
"""Build the PageIndex ``IndexConfig`` for local indexing.
158158
159159
Forwards the KB's ``concurrency`` setting to PageIndex, which caps how many
@@ -162,6 +162,15 @@ def _build_index_config(config: dict[str, Any]) -> IndexConfig:
162162
installed PageIndex's ``IndexConfig`` declares the field, so OpenKB keeps
163163
working against a pinned PageIndex that predates it (``IndexConfig``
164164
forbids unknown kwargs).
165+
166+
``bundle``'s ``api_key``/``base_url`` (the same credentials ``compiler.py``'s
167+
own LLM calls use — see :func:`openkb.config.resolve_credential_bundle`) are
168+
forwarded as PageIndex's own per-call ``llm_params`` (see #219): without
169+
this, PageIndex's internal indexing calls (TOC/tree/summary generation)
170+
fall back to LiteLLM's default provider-key/env-var lookup, which doesn't
171+
know about a KB's custom ``LLM_API_KEY``/gateway ``base_url``. Guarded by
172+
the same ``model_fields`` check as ``max_concurrency`` above, so it
173+
degrades gracefully against an older pinned PageIndex.
165174
"""
166175
kwargs: dict[str, Any] = {
167176
"if_add_node_text": True,
@@ -177,9 +186,48 @@ def _build_index_config(config: dict[str, Any]) -> IndexConfig:
177186
"config: 'concurrency' is set but the installed PageIndex "
178187
"version does not support it yet — ignoring it."
179188
)
189+
if bundle is not None:
190+
llm_params = {
191+
key: value
192+
for key, value in {"api_key": bundle.api_key, "base_url": bundle.base_url}.items()
193+
if value
194+
}
195+
if llm_params:
196+
if "llm_params" in IndexConfig.model_fields:
197+
kwargs["llm_params"] = llm_params
198+
else:
199+
logger.warning(
200+
"config: a custom LLM_API_KEY/base_url is set but the installed "
201+
"PageIndex version doesn't support forwarding it (llm_params) yet "
202+
"— PageIndex's own LLM calls will use their default credential "
203+
"lookup instead."
204+
)
180205
return IndexConfig(**kwargs)
181206

182207

208+
def _close_pageindex_client(client: PageIndexClient) -> None:
209+
"""Best-effort close of ``client``'s local SQLite connection(s).
210+
211+
The pinned ``pageindex`` version has no public ``close()``/context-manager
212+
API on ``PageIndexClient``/``Collection`` — only on the low-level
213+
``SQLiteStorage`` its local backend holds internally — so this reaches into
214+
that private attribute directly. In *cloud* mode (``client`` built with a
215+
``PAGEINDEX_API_KEY``) there is no local backend/storage at all, so this is
216+
a silent no-op. Never raises: called from both the success and the failure
217+
path of :func:`index_long_document`, and closing must never mask a real
218+
indexing error. Without this, a subsequent mutation rollback can't
219+
unlink/rename ``pageindex.db`` on Windows while this process still holds it
220+
open (#249).
221+
"""
222+
storage = getattr(getattr(client, "_backend", None), "_storage", None)
223+
close = getattr(storage, "close", None)
224+
if callable(close):
225+
try:
226+
close()
227+
except Exception:
228+
logger.debug("Failed to close PageIndex's local storage", exc_info=True)
229+
230+
183231
def index_long_document(pdf_path: Path, kb_dir: Path, doc_name: str | None = None) -> IndexResult:
184232
"""Index a long PDF document using PageIndex and write wiki pages.
185233
@@ -192,106 +240,115 @@ def index_long_document(pdf_path: Path, kb_dir: Path, doc_name: str | None = Non
192240

193241
model: str = config.get("model", "gpt-5.4")
194242
pageindex_api_key = os.environ.get("PAGEINDEX_API_KEY", "")
243+
bundle = resolve_credential_bundle(kb_dir)
195244

196-
index_config = _build_index_config(config)
245+
index_config = _build_index_config(config, bundle)
197246

198247
client = PageIndexClient(
199248
api_key=pageindex_api_key or None,
200249
model=model,
201250
storage_path=str(openkb_dir),
202251
index_config=index_config,
203252
)
204-
col = client.collection()
205-
206-
# Add PDF (retry up to 3 times — PageIndex TOC accuracy is stochastic)
207-
max_retries = 3
208-
doc_id = None
209-
for attempt in range(1, max_retries + 1):
210-
try:
211-
doc_id = col.add(str(pdf_path))
212-
logger.info(
213-
"PageIndex added %s → doc_id=%s (attempt %d)", pdf_path.name, doc_id, attempt
214-
)
215-
break
216-
except Exception as exc:
217-
logger.warning(
218-
"PageIndex attempt %d/%d failed for %s: %s",
219-
attempt,
220-
max_retries,
221-
pdf_path.name,
222-
exc,
223-
)
224-
if attempt == max_retries:
225-
raise RuntimeError(
226-
f"Failed to index {pdf_path.name} after {max_retries} attempts: {exc}"
227-
) from exc
228-
229-
# The PageIndex blob for doc_id is now durably on disk. The add mutation no
230-
# longer eagerly snapshots .openkb/files — it registers the new blob via
231-
# snapshot.track_new() only on a successful return — so if any step below
232-
# fails, delete the document we just added. Otherwise the blob leaks as an
233-
# orphan that pageindex.db (rolled back by the snapshot) no longer refs and
234-
# no reaper reclaims.
253+
# Closed in `finally` (both success and failure) so a subsequent mutation
254+
# rollback can always unlink/rename pageindex.db on Windows — see #249.
235255
try:
236-
# Fetch complete document (metadata + structure + text)
237-
doc = col.get_document(doc_id, include_text=True)
238-
indexed_doc_name: str = doc.get("doc_name", pdf_path.stem)
239-
description: str = doc.get("doc_description", "")
240-
structure: list = doc.get("structure", [])
241-
242-
# Debug: print doc keys and page_count to diagnose get_page_content range
243-
logger.info("Doc keys: %s", list(doc.keys()))
244-
logger.info("page_count from doc: %s", doc.get("page_count", "NOT PRESENT"))
245-
246-
tree = {
247-
"doc_name": indexed_doc_name,
248-
"doc_description": description,
249-
"structure": structure,
250-
}
251-
252-
# Write wiki/sources/ — per-page content
253-
sources_dir = kb_dir / "wiki" / "sources"
254-
sources_dir.mkdir(parents=True, exist_ok=True)
255-
images_dir = sources_dir / "images" / source_name
256+
col = client.collection()
256257

257-
all_pages: list[dict[str, Any]] = []
258-
if pageindex_api_key:
259-
# Cloud mode: fetch OCR'd markdown from PageIndex. get_page_content
260-
# requires a page range, so pass "1-N".
261-
page_count = _get_pdf_page_count(pdf_path)
258+
# Add PDF (retry up to 3 times — PageIndex TOC accuracy is stochastic)
259+
max_retries = 3
260+
doc_id = None
261+
for attempt in range(1, max_retries + 1):
262262
try:
263-
all_pages = _normalize_page_content(col.get_page_content(doc_id, f"1-{page_count}"))
263+
doc_id = col.add(str(pdf_path))
264+
logger.info(
265+
"PageIndex added %s → doc_id=%s (attempt %d)", pdf_path.name, doc_id, attempt
266+
)
267+
break
264268
except Exception as exc:
265-
logger.warning("Cloud get_page_content failed for %s: %s", pdf_path.name, exc)
266-
267-
if not all_pages:
268-
if pageindex_api_key:
269269
logger.warning(
270-
"Cloud returned no pages for %s; falling back to local pymupdf", pdf_path.name
270+
"PageIndex attempt %d/%d failed for %s: %s",
271+
attempt,
272+
max_retries,
273+
pdf_path.name,
274+
exc,
275+
)
276+
if attempt == max_retries:
277+
raise RuntimeError(
278+
f"Failed to index {pdf_path.name} after {max_retries} attempts: {exc}"
279+
) from exc
280+
281+
# The PageIndex blob for doc_id is now durably on disk. The add mutation no
282+
# longer eagerly snapshots .openkb/files — it registers the new blob via
283+
# snapshot.track_new() only on a successful return — so if any step below
284+
# fails, delete the document we just added. Otherwise the blob leaks as an
285+
# orphan that pageindex.db (rolled back by the snapshot) no longer refs and
286+
# no reaper reclaims.
287+
try:
288+
# Fetch complete document (metadata + structure + text)
289+
doc = col.get_document(doc_id, include_text=True)
290+
indexed_doc_name: str = doc.get("doc_name", pdf_path.stem)
291+
description: str = doc.get("doc_description", "")
292+
structure: list = doc.get("structure", [])
293+
294+
# Debug: print doc keys and page_count to diagnose get_page_content range
295+
logger.info("Doc keys: %s", list(doc.keys()))
296+
logger.info("page_count from doc: %s", doc.get("page_count", "NOT PRESENT"))
297+
298+
tree = {
299+
"doc_name": indexed_doc_name,
300+
"doc_description": description,
301+
"structure": structure,
302+
}
303+
304+
# Write wiki/sources/ — per-page content
305+
sources_dir = kb_dir / "wiki" / "sources"
306+
sources_dir.mkdir(parents=True, exist_ok=True)
307+
images_dir = sources_dir / "images" / source_name
308+
309+
all_pages: list[dict[str, Any]] = []
310+
if pageindex_api_key:
311+
# Cloud mode: fetch OCR'd markdown from PageIndex. get_page_content
312+
# requires a page range, so pass "1-N".
313+
page_count = _get_pdf_page_count(pdf_path)
314+
try:
315+
all_pages = _normalize_page_content(
316+
col.get_page_content(doc_id, f"1-{page_count}")
317+
)
318+
except Exception as exc:
319+
logger.warning("Cloud get_page_content failed for %s: %s", pdf_path.name, exc)
320+
321+
if not all_pages:
322+
if pageindex_api_key:
323+
logger.warning(
324+
"Cloud returned no pages for %s; falling back to local pymupdf",
325+
pdf_path.name,
326+
)
327+
all_pages = _normalize_page_content(
328+
_convert_pdf_to_pages(pdf_path, source_name, images_dir)
271329
)
272-
all_pages = _normalize_page_content(
273-
_convert_pdf_to_pages(pdf_path, source_name, images_dir)
274-
)
275330

276-
if not all_pages:
277-
raise RuntimeError(f"No page content extracted for {pdf_path.name}")
331+
if not all_pages:
332+
raise RuntimeError(f"No page content extracted for {pdf_path.name}")
278333

279-
_write_long_doc_artifacts(
280-
tree, all_pages, source_name, doc_id, kb_dir, description=description
281-
)
282-
return IndexResult(doc_id=doc_id, description=description, tree=tree)
283-
except BaseException:
284-
# Best-effort: remove the blob this add created. A failure here (e.g. a
285-
# second interrupt) only means the blob may stay orphaned — the original
286-
# error still propagates so the caller (mutation coordinator) rolls back
287-
# everything else it snapshotted.
288-
try:
289-
col.delete_document(doc_id)
290-
except Exception:
291-
logger.warning(
292-
"PageIndex cleanup of %s failed after error; blob may be orphaned", doc_id
334+
_write_long_doc_artifacts(
335+
tree, all_pages, source_name, doc_id, kb_dir, description=description
293336
)
294-
raise
337+
return IndexResult(doc_id=doc_id, description=description, tree=tree)
338+
except BaseException:
339+
# Best-effort: remove the blob this add created. A failure here (e.g. a
340+
# second interrupt) only means the blob may stay orphaned — the original
341+
# error still propagates so the caller (mutation coordinator) rolls back
342+
# everything else it snapshotted.
343+
try:
344+
col.delete_document(doc_id)
345+
except Exception:
346+
logger.warning(
347+
"PageIndex cleanup of %s failed after error; blob may be orphaned", doc_id
348+
)
349+
raise
350+
finally:
351+
_close_pageindex_client(client)
295352

296353

297354
# PageIndex's get_page_content rejects a single page range covering more than

0 commit comments

Comments
 (0)