Files

324 lines
12 KiB
Python

"""Docling document conversion for NAS files."""
from __future__ import annotations
import json
import logging
import os
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
import psycopg2
log = logging.getLogger("doc-ingest.docling")
_executor = ThreadPoolExecutor(max_workers=1)
DB_HOST = os.getenv("DB_HOST", "10.4.7.18")
DB_USER = os.getenv("DB_USER", "aissa")
DB_PASSWORD = os.getenv("DB_PASSWORD", "Foodlinkk#2026")
DB_NAME = os.getenv("DB_NAME", "foodlinkk")
SUPPORTED_SUFFIXES = {
".pdf", ".docx", ".pptx", ".xlsx", ".html", ".htm", ".md", ".txt",
".png", ".jpg", ".jpeg", ".webp", ".tiff", ".tif", ".bmp", ".gif",
}
EXPORT_FORMATS = ("markdown", "html", "json", "doctags", "text", "yaml")
def _db():
return psycopg2.connect(host=DB_HOST, user=DB_USER, password=DB_PASSWORD, dbname=DB_NAME)
def file_signature(path: Path) -> str:
st = path.stat()
return f"{int(st.st_mtime)}_{st.st_size}"
def capabilities() -> dict[str, Any]:
try:
import docling # noqa: F401
installed = True
version = getattr(docling, "__version__", "unknown")
except Exception:
installed = False
version = None
return {
"installed": installed,
"version": version,
"supported_suffixes": sorted(SUPPORTED_SUFFIXES),
"export_formats": list(EXPORT_FORMATS),
"pipeline_options": [
{"id": "ocr", "label": "OCR", "default": True},
{"id": "tables", "label": "Tabelstructuur", "default": True},
{"id": "page_images", "label": "Pagina-afbeeldingen", "default": False},
{"id": "picture_images", "label": "Figuur-afbeeldingen", "default": True},
{"id": "force_full_page_ocr", "label": "Volledige pagina OCR", "default": False},
],
}
def _build_converter(options: dict[str, Any]):
from docling.document_converter import DocumentConverter, PdfFormatOption
from docling.datamodel.base_models import InputFormat
from docling.datamodel.pipeline_options import PdfPipelineOptions
pipeline_options = PdfPipelineOptions()
pipeline_options.do_ocr = bool(options.get("ocr", True))
pipeline_options.do_table_structure = bool(options.get("tables", True))
pipeline_options.generate_page_images = bool(options.get("page_images", False))
pipeline_options.generate_picture_images = bool(options.get("picture_images", True))
if hasattr(pipeline_options, "force_full_page_ocr"):
pipeline_options.force_full_page_ocr = bool(options.get("force_full_page_ocr", False))
return DocumentConverter(
format_options={
InputFormat.PDF: PdfFormatOption(pipeline_options=pipeline_options),
}
)
def _export_document(doc: Any, fmt: str) -> str | dict | list:
if fmt == "markdown":
return doc.export_to_markdown()
if fmt == "html":
return doc.export_to_html()
if fmt == "json":
if hasattr(doc, "export_to_dict"):
return doc.export_to_dict()
if hasattr(doc, "model_dump"):
return doc.model_dump()
return json.loads(doc.model_dump_json()) if hasattr(doc, "model_dump_json") else {}
if fmt == "doctags":
if hasattr(doc, "export_to_doctags"):
return doc.export_to_doctags()
if hasattr(doc, "export_to_document_tokens"):
return doc.export_to_document_tokens()
return ""
if fmt == "text":
if hasattr(doc, "export_to_text"):
return doc.export_to_text()
return doc.export_to_markdown()
if fmt == "yaml":
if hasattr(doc, "export_to_yaml"):
return doc.export_to_yaml()
return ""
raise ValueError(f"Unknown export format: {fmt}")
def _extract_tables(doc: Any) -> list[dict[str, Any]]:
tables: list[dict[str, Any]] = []
for idx, table in enumerate(getattr(doc, "tables", []) or []):
row: dict[str, Any] = {"index": idx}
try:
if hasattr(table, "export_to_markdown"):
row["markdown"] = table.export_to_markdown()
if hasattr(table, "export_to_dataframe"):
df = table.export_to_dataframe()
row["csv"] = df.to_csv(index=False)
row["rows"] = df.values.tolist()
row["columns"] = list(df.columns)
except Exception as exc:
row["error"] = str(exc)[:200]
tables.append(row)
return tables
def _extract_pictures(doc: Any) -> list[dict[str, Any]]:
pictures: list[dict[str, Any]] = []
for idx, pic in enumerate(getattr(doc, "pictures", []) or []):
item: dict[str, Any] = {"index": idx}
prov = getattr(pic, "prov", None) or []
if prov:
item["page"] = getattr(prov[0], "page_no", None)
if hasattr(pic, "caption_text"):
item["caption"] = pic.caption_text()
pictures.append(item)
return pictures
def convert_path(path: Path, options: dict[str, Any]) -> dict[str, Any]:
if not path.is_file():
raise FileNotFoundError(str(path))
if path.suffix.lower() not in SUPPORTED_SUFFIXES:
raise ValueError(f"Unsupported type: {path.suffix}")
formats = options.get("formats") or ["markdown", "json"]
formats = [f for f in formats if f in EXPORT_FORMATS]
if not formats:
formats = ["markdown"]
converter = _build_converter(options)
result = converter.convert(str(path))
doc = result.document
exports: dict[str, Any] = {}
for fmt in formats:
try:
exports[fmt] = _export_document(doc, fmt)
except Exception as exc:
exports[fmt] = {"error": str(exc)[:300]}
meta = {}
if hasattr(doc, "name"):
meta["name"] = doc.name
if hasattr(result, "input") and result.input:
meta["input"] = str(result.input.file) if hasattr(result.input, "file") else str(result.input)
return {
"exports": exports,
"tables": _extract_tables(doc),
"pictures": _extract_pictures(doc),
"page_count": len(getattr(doc, "pages", []) or []),
"metadata": meta,
}
def _upsert_result(
storage_path: str,
file_sig: str,
status: str,
options: dict,
payload: dict | None = None,
error: str | None = None,
) -> int:
payload = payload or {}
conn = _db()
conn.autocommit = True
with conn.cursor() as cur:
cur.execute(
"""
INSERT INTO docling_results
(storage_path, file_sig, status, options, exports, tables_data, pictures,
page_count, metadata, error_text, completed_at)
VALUES (%s, %s, %s, %s::jsonb, %s::jsonb, %s::jsonb, %s::jsonb, %s, %s::jsonb, %s, %s)
ON CONFLICT (storage_path, file_sig) DO UPDATE SET
status = EXCLUDED.status,
options = EXCLUDED.options,
exports = EXCLUDED.exports,
tables_data = EXCLUDED.tables_data,
pictures = EXCLUDED.pictures,
page_count = EXCLUDED.page_count,
metadata = EXCLUDED.metadata,
error_text = EXCLUDED.error_text,
completed_at = EXCLUDED.completed_at
RETURNING id
""",
(
storage_path,
file_sig,
status,
json.dumps(options),
json.dumps(payload.get("exports", {})),
json.dumps(payload.get("tables", [])),
json.dumps(payload.get("pictures", [])),
payload.get("page_count", 0),
json.dumps(payload.get("metadata", {})),
error,
datetime.now(timezone.utc) if status in ("done", "failed") else None,
),
)
row_id = cur.fetchone()[0]
conn.close()
return row_id
def get_result(storage_path: str, file_sig: str | None = None) -> dict | None:
conn = _db()
with conn.cursor() as cur:
if file_sig:
cur.execute(
"""
SELECT id, storage_path, file_sig, status, options, exports, tables_data, pictures,
page_count, metadata, error_text, created_at, completed_at
FROM docling_results WHERE storage_path = %s AND file_sig = %s
""",
(storage_path, file_sig),
)
else:
cur.execute(
"""
SELECT id, storage_path, file_sig, status, options, exports, tables_data, pictures,
page_count, metadata, error_text, created_at, completed_at
FROM docling_results WHERE storage_path = %s
ORDER BY completed_at DESC NULLS LAST, created_at DESC LIMIT 1
""",
(storage_path,),
)
row = cur.fetchone()
conn.close()
if not row:
return None
cols = [
"id", "storage_path", "file_sig", "status", "options", "exports", "tables_data",
"pictures", "page_count", "metadata", "error_text", "created_at", "completed_at",
]
out = dict(zip(cols, row))
for key in ("created_at", "completed_at"):
if out.get(key):
out[key] = out[key].isoformat()
return out
def list_history(limit: int = 40) -> list[dict[str, Any]]:
limit = max(1, min(limit, 100))
conn = _db()
with conn.cursor() as cur:
cur.execute(
"""
SELECT id, storage_path, file_sig, status, page_count,
jsonb_array_length(COALESCE(tables_data, '[]'::jsonb)) AS table_count,
jsonb_array_length(COALESCE(pictures, '[]'::jsonb)) AS picture_count,
completed_at, error_text,
LEFT(COALESCE(exports->>'markdown', ''), 120) AS md_preview
FROM docling_results
ORDER BY completed_at DESC NULLS LAST, created_at DESC
LIMIT %s
""",
(limit,),
)
rows = cur.fetchall()
conn.close()
cols = [
"id", "storage_path", "file_sig", "status", "page_count", "table_count",
"picture_count", "completed_at", "error_text", "md_preview",
]
out = []
for row in rows:
item = dict(zip(cols, row))
if item.get("completed_at"):
item["completed_at"] = item["completed_at"].isoformat()
out.append(item)
return out
def convert_nas_file(nas_root: Path, rel_path: str, options: dict[str, Any]) -> dict[str, Any]:
path = (nas_root / rel_path).resolve()
if not str(path).startswith(str(nas_root.resolve())):
raise ValueError("Path outside NAS root")
sig = file_signature(path)
_upsert_result(rel_path, sig, "running", options)
try:
payload = convert_path(path, options)
job_id = _upsert_result(rel_path, sig, "done", options, payload)
return {"ok": True, "id": job_id, "storage_path": rel_path, "file_sig": sig, **payload}
except Exception as exc:
job_id = _upsert_result(rel_path, sig, "failed", options, error=str(exc)[:500])
return {"ok": False, "id": job_id, "storage_path": rel_path, "error": str(exc)[:500]}
def convert_nas_batch(nas_root: Path, rel_paths: list[str], options: dict[str, Any]) -> dict[str, Any]:
results = []
ok_count = 0
for rel in rel_paths:
try:
out = convert_nas_file(nas_root, rel, options)
if out.get("ok"):
ok_count += 1
results.append(out)
except Exception as exc:
results.append({"ok": False, "storage_path": rel, "error": str(exc)[:300]})
return {"ok": True, "converted": ok_count, "total": len(rel_paths), "results": results}