mirror of
https://github.com/paperless-ngx/paperless-ngx.git
synced 2026-07-30 07:44:54 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7d9ba34582 | ||
|
|
91c3d9caab | ||
|
|
668fa77428 | ||
|
|
5bd72014a6 | ||
|
|
bbb9c86ba4 |
@@ -948,11 +948,10 @@ for display in the web interface.
|
||||
|
||||
!!! note
|
||||
|
||||
The **remote OCR parser** (Azure AI) also honors this setting: when
|
||||
no archive is requested (`never`, or `auto` with a born-digital PDF),
|
||||
the remote engine is skipped entirely and locally-extracted text is
|
||||
used instead, avoiding an unnecessary API call and a duplicate text
|
||||
layer.
|
||||
The **remote OCR parser** (Azure AI) always produces a searchable
|
||||
PDF and stores it as the archive copy, regardless of this setting.
|
||||
`ARCHIVE_FILE_GENERATION=never` has no effect when the remote
|
||||
parser handles a document.
|
||||
|
||||
#### [`PAPERLESS_OCR_CLEAN=<mode>`](#PAPERLESS_OCR_CLEAN) {#PAPERLESS_OCR_CLEAN}
|
||||
|
||||
|
||||
@@ -187,11 +187,10 @@ PAPERLESS_ARCHIVE_FILE_GENERATION=auto
|
||||
|
||||
### Remote OCR parser
|
||||
|
||||
If you use the **remote OCR parser** (Azure AI), `ARCHIVE_FILE_GENERATION` is
|
||||
honored the same way as for the local engine: when no archive is requested
|
||||
(`never`, or `auto` with a born-digital PDF), the remote engine is skipped
|
||||
entirely and locally-extracted text is used instead, avoiding an unnecessary
|
||||
API call and a duplicate text layer.
|
||||
If you use the **remote OCR parser** (Azure AI), note that it always produces a
|
||||
searchable PDF and stores it as the archive copy. `ARCHIVE_FILE_GENERATION=never`
|
||||
has no effect for documents handled by the remote parser - the archive is produced
|
||||
unconditionally by the remote engine.
|
||||
|
||||
## Search Index (Whoosh -> Tantivy)
|
||||
|
||||
|
||||
@@ -129,13 +129,25 @@ describe('PngxPdfViewerComponent', () => {
|
||||
;(component as any).applyScale()
|
||||
expect(viewer.currentScaleValue).toBe(PdfZoomScale.PageFit)
|
||||
expect(viewer.currentScale).toBe(2)
|
||||
})
|
||||
|
||||
it('does not reapply scale for page-only changes', async () => {
|
||||
await initComponent()
|
||||
|
||||
const pdf = (component as any).pdf as { numPages: number }
|
||||
pdf.numPages = 3
|
||||
const viewer = (component as any).pdfViewer as PDFViewer
|
||||
viewer.setDocument(pdf)
|
||||
const applyScaleSpy = jest.spyOn(component as any, 'applyScale')
|
||||
component.page = 2
|
||||
;(component as any).lastViewerPage = 2
|
||||
;(component as any).applyViewerState()
|
||||
|
||||
component.ngOnChanges({
|
||||
page: new SimpleChange(1, 2, false),
|
||||
})
|
||||
|
||||
expect(viewer.currentPageNumber).toBe(2)
|
||||
expect((component as any).lastViewerPage).toBeUndefined()
|
||||
expect(applyScaleSpy).toHaveBeenCalled()
|
||||
expect(applyScaleSpy).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('does not reset the viewer when it is already on the requested page', async () => {
|
||||
|
||||
@@ -116,7 +116,10 @@ export class PngxPdfViewerComponent
|
||||
changes['zoomScale'] ||
|
||||
changes['rotation']
|
||||
) {
|
||||
this.applyViewerState()
|
||||
// Prevent loop with page / scale application see https://github.com/paperless-ngx/paperless-ngx/issues/13404
|
||||
this.applyViewerState(
|
||||
!!(changes['zoom'] || changes['zoomScale'] || changes['rotation'])
|
||||
)
|
||||
}
|
||||
|
||||
if (changes['searchQuery']) {
|
||||
@@ -240,7 +243,7 @@ export class PngxPdfViewerComponent
|
||||
}
|
||||
}
|
||||
|
||||
private applyViewerState(): void {
|
||||
private applyViewerState(applyScale = true): void {
|
||||
if (!this.pdfViewer) {
|
||||
return
|
||||
}
|
||||
@@ -264,7 +267,7 @@ export class PngxPdfViewerComponent
|
||||
if (this.page === this.lastViewerPage) {
|
||||
this.lastViewerPage = undefined
|
||||
}
|
||||
if (hasPages) {
|
||||
if (hasPages && applyScale) {
|
||||
this.applyScale()
|
||||
}
|
||||
this.dispatchFindIfReady()
|
||||
|
||||
@@ -36,6 +36,9 @@ def send_email(
|
||||
|
||||
TODO: re-evaluate this pending https://code.djangoproject.com/ticket/35581 / https://github.com/django/django/pull/18966
|
||||
"""
|
||||
if "\r" in subject or "\n" in subject:
|
||||
subject = " ".join(line.strip(" \t") for line in subject.splitlines())
|
||||
|
||||
email = EmailMessage(
|
||||
subject=subject,
|
||||
body=body,
|
||||
|
||||
@@ -386,10 +386,19 @@ class Command(CryptMixin, PaperlessCommand):
|
||||
raise DeserializationError(
|
||||
f"{model.__name__} has no updatable fields; PK-only models are not supported by the importer",
|
||||
)
|
||||
# MySQL/MariaDB support upserts via ON DUPLICATE KEY UPDATE but,
|
||||
# unlike PostgreSQL/SQLite, cannot target a specific unique field
|
||||
# for the conflict -- passing unique_fields there raises
|
||||
# NotSupportedError.
|
||||
unique_fields = (
|
||||
[model._meta.pk.attname]
|
||||
if connection.features.supports_update_conflicts_with_target
|
||||
else None
|
||||
)
|
||||
model.objects.bulk_create( # type: ignore[attr-defined]
|
||||
instances,
|
||||
update_conflicts=True,
|
||||
unique_fields=[model._meta.pk.attname],
|
||||
unique_fields=unique_fields,
|
||||
update_fields=update_fields,
|
||||
)
|
||||
loaded_models.add(model)
|
||||
|
||||
@@ -75,7 +75,7 @@ class TestEmail(DirectoriesMixin, SampleDirMixin, APITestCase):
|
||||
{
|
||||
"documents": [self.doc1.pk, self.doc2.pk],
|
||||
"addresses": "hello@paperless-ngx.com,test@example.com",
|
||||
"subject": "Bulk email test",
|
||||
"subject": "Bulk email\n test",
|
||||
"message": "Here are your documents",
|
||||
},
|
||||
),
|
||||
|
||||
@@ -3,9 +3,7 @@ Built-in remote-OCR document parser.
|
||||
|
||||
Handles documents by sending them to a configured remote OCR engine
|
||||
(currently Azure AI Vision / Document Intelligence) and retrieving both
|
||||
the extracted text and a searchable PDF with an embedded text layer. For
|
||||
born-digital PDFs that need no archive copy, the remote call is skipped
|
||||
entirely in favor of locally-extracted text (see ``RemoteDocumentParser.parse``).
|
||||
the extracted text and a searchable PDF with an embedded text layer.
|
||||
|
||||
When no engine is configured, ``score()`` returns ``None`` so the parser
|
||||
is effectively invisible to the registry — the tesseract parser handles
|
||||
@@ -23,8 +21,6 @@ from typing import Self
|
||||
|
||||
from django.conf import settings
|
||||
|
||||
from paperless.parsers.utils import extract_pdf_text
|
||||
from paperless.parsers.utils import post_process_text
|
||||
from paperless.version import __full_version_str__
|
||||
|
||||
if TYPE_CHECKING:
|
||||
@@ -73,11 +69,8 @@ class RemoteDocumentParser:
|
||||
"""Parse documents via a remote OCR API (currently Azure AI Vision).
|
||||
|
||||
This parser sends documents to a remote engine that returns both
|
||||
extracted text and a searchable PDF with an embedded text layer,
|
||||
except when ``parse()`` is called with ``produce_archive=False`` for
|
||||
a PDF, in which case the remote call is skipped and only locally
|
||||
extracted text is returned (no archive). It does not depend on
|
||||
Tesseract or ocrmypdf.
|
||||
extracted text and a searchable PDF with an embedded text layer.
|
||||
It does not depend on Tesseract or ocrmypdf.
|
||||
|
||||
Class attributes
|
||||
----------------
|
||||
@@ -166,11 +159,8 @@ class RemoteDocumentParser:
|
||||
Returns
|
||||
-------
|
||||
bool
|
||||
Always True — the remote engine is capable of returning a PDF
|
||||
with an embedded text layer to serve as the archive copy.
|
||||
Whether it actually does so for a given document depends on
|
||||
``produce_archive`` passed to :meth:`parse` (see there for when
|
||||
the remote engine call, and thus archive generation, is skipped).
|
||||
Always True — the remote engine always returns a PDF with an
|
||||
embedded text layer that serves as the archive copy.
|
||||
"""
|
||||
return True
|
||||
|
||||
@@ -227,12 +217,6 @@ class RemoteDocumentParser:
|
||||
) -> None:
|
||||
"""Send the document to the remote engine and store results.
|
||||
|
||||
When *produce_archive* is False for a PDF, the caller (via
|
||||
``documents.consumer.should_produce_archive``) has already determined
|
||||
that the document is born-digital and needs no archive — skip the
|
||||
remote engine entirely rather than re-OCRing it and creating a
|
||||
duplicate text layer.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
document_path:
|
||||
@@ -240,8 +224,8 @@ class RemoteDocumentParser:
|
||||
mime_type:
|
||||
Detected MIME type of the document.
|
||||
produce_archive:
|
||||
Whether an archive copy is wanted. For PDFs, False skips the
|
||||
remote engine and uses locally-extracted text instead.
|
||||
Ignored — the remote engine always returns a searchable PDF,
|
||||
which is stored as the archive copy regardless of this flag.
|
||||
"""
|
||||
config = RemoteEngineConfig(
|
||||
engine=settings.REMOTE_OCR_ENGINE,
|
||||
@@ -256,16 +240,6 @@ class RemoteDocumentParser:
|
||||
self._text = ""
|
||||
return
|
||||
|
||||
if not produce_archive and mime_type == "application/pdf":
|
||||
logger.debug(
|
||||
"Remote OCR: skipped — no archive requested, "
|
||||
"using locally-extracted text",
|
||||
)
|
||||
self._text = (
|
||||
post_process_text(extract_pdf_text(document_path, log=logger)) or ""
|
||||
)
|
||||
return
|
||||
|
||||
if config.engine == "azureai":
|
||||
self._text = self._azure_ai_vision_parse(document_path, config)
|
||||
|
||||
|
||||
@@ -3,6 +3,7 @@ from __future__ import annotations
|
||||
import importlib.resources
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
import shutil
|
||||
import tempfile
|
||||
from pathlib import Path
|
||||
@@ -24,9 +25,9 @@ from paperless.config import OcrConfig
|
||||
from paperless.models import CleanChoices
|
||||
from paperless.models import ModeChoices
|
||||
from paperless.models import OutputTypeChoices
|
||||
from paperless.parsers.utils import PDF_TEXT_MIN_LENGTH
|
||||
from paperless.parsers.utils import extract_pdf_text
|
||||
from paperless.parsers.utils import pdf_has_digital_text
|
||||
from paperless.parsers.utils import post_process_text
|
||||
from paperless.parsers.utils import is_tagged_pdf
|
||||
from paperless.parsers.utils import read_file_handle_unicode_errors
|
||||
from paperless.version import __full_version_str__
|
||||
|
||||
@@ -509,10 +510,10 @@ class RasterisedDocumentParser:
|
||||
|
||||
if mime_type == "application/pdf":
|
||||
text_original = self.extract_text(None, document_path)
|
||||
original_has_text = pdf_has_digital_text(
|
||||
document_path,
|
||||
text_original,
|
||||
log=self.log,
|
||||
has_text = text_original is not None and len(text_original) > 0
|
||||
original_has_text = has_text and (
|
||||
is_tagged_pdf(document_path, log=self.log)
|
||||
or len(text_original) > PDF_TEXT_MIN_LENGTH
|
||||
)
|
||||
else:
|
||||
text_original = None
|
||||
@@ -657,3 +658,17 @@ class RasterisedDocumentParser:
|
||||
f"No text was found in {document_path}, the content will be empty.",
|
||||
)
|
||||
self.text = ""
|
||||
|
||||
|
||||
def post_process_text(text: str | None) -> str | None:
|
||||
if not text:
|
||||
return None
|
||||
|
||||
collapsed_spaces = re.sub(r"([^\S\r\n]+)", " ", text)
|
||||
no_leading_whitespace = re.sub(r"([\n\r]+)([^\S\n\r]+)", "\\1", collapsed_spaces)
|
||||
no_trailing_whitespace = re.sub(r"([^\S\n\r]+)$", "", no_leading_whitespace)
|
||||
|
||||
# TODO: this needs a rework
|
||||
# replace \0 prevents issues with saving to postgres.
|
||||
# text may contain \0 when this character is present in PDF files.
|
||||
return no_trailing_whitespace.strip().replace("\0", " ")
|
||||
|
||||
@@ -65,63 +65,6 @@ def is_tagged_pdf(
|
||||
return False
|
||||
|
||||
|
||||
def pdf_has_digital_text(
|
||||
path: Path,
|
||||
text: str | None,
|
||||
log: logging.Logger | None = None,
|
||||
) -> bool:
|
||||
"""Return True if a PDF already has a usable, born-digital text layer.
|
||||
|
||||
Combines the tagged-PDF check with an extracted-text length check.
|
||||
Shared by the tesseract and remote OCR parsers to decide whether
|
||||
OCR_MODE=auto/off should skip (re-)OCRing a document.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
path:
|
||||
Absolute path to the PDF file.
|
||||
text:
|
||||
Text already extracted from the PDF (e.g. via ``extract_pdf_text``),
|
||||
or ``None``.
|
||||
log:
|
||||
Logger for warnings. Falls back to the module-level logger when omitted.
|
||||
|
||||
Returns
|
||||
-------
|
||||
bool
|
||||
``True`` when the document already contains a text layer.
|
||||
"""
|
||||
return is_tagged_pdf(path, log=log) or (
|
||||
text is not None and len(text) > PDF_TEXT_MIN_LENGTH
|
||||
)
|
||||
|
||||
|
||||
def post_process_text(text: str | None) -> str | None:
|
||||
"""Normalise whitespace in extracted OCR/PDF text and strip NUL bytes.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
text:
|
||||
Raw extracted text, or ``None``.
|
||||
|
||||
Returns
|
||||
-------
|
||||
str | None
|
||||
Cleaned text, or ``None`` when *text* is falsy.
|
||||
"""
|
||||
if not text:
|
||||
return None
|
||||
|
||||
collapsed_spaces = re.sub(r"([^\S\r\n]+)", " ", text)
|
||||
no_leading_whitespace = re.sub(r"([\n\r]+)([^\S\n\r]+)", "\\1", collapsed_spaces)
|
||||
no_trailing_whitespace = re.sub(r"([^\S\n\r]+)$", "", no_leading_whitespace)
|
||||
|
||||
# TODO: this needs a rework
|
||||
# replace \0 prevents issues with saving to postgres.
|
||||
# text may contain \0 when this character is present in PDF files.
|
||||
return no_trailing_whitespace.strip().replace("\0", " ")
|
||||
|
||||
|
||||
def extract_pdf_text(
|
||||
path: Path,
|
||||
log: logging.Logger | None = None,
|
||||
|
||||
@@ -336,117 +336,6 @@ class TestRemoteParserParse:
|
||||
assert remote_parser.get_date() is None
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# parse() — produce_archive=False skips the remote engine (PDFs only)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TestRemoteParserSkipsWhenNoArchiveWanted:
|
||||
"""When the caller has already decided no archive is needed for a PDF
|
||||
(documents.consumer.should_produce_archive), the remote engine call is
|
||||
skipped entirely in favor of locally-extracted text.
|
||||
"""
|
||||
|
||||
def test_pdf_skips_azure_when_no_archive_requested(
|
||||
self,
|
||||
remote_parser: RemoteDocumentParser,
|
||||
simple_digital_pdf_file: Path,
|
||||
azure_client: Mock,
|
||||
) -> None:
|
||||
"""
|
||||
GIVEN: produce_archive=False for a PDF
|
||||
WHEN: parse() is called
|
||||
THEN: Azure is never invoked, no archive is produced, and text
|
||||
comes from local pdftotext extraction
|
||||
"""
|
||||
remote_parser.parse(
|
||||
simple_digital_pdf_file,
|
||||
"application/pdf",
|
||||
produce_archive=False,
|
||||
)
|
||||
|
||||
azure_client.begin_analyze_document.assert_not_called()
|
||||
assert remote_parser.get_archive_path() is None
|
||||
assert remote_parser.get_text() != ""
|
||||
|
||||
def test_pdf_no_archive_requested_text_matches_local_extraction(
|
||||
self,
|
||||
remote_parser: RemoteDocumentParser,
|
||||
simple_digital_pdf_file: Path,
|
||||
azure_client: Mock,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""
|
||||
GIVEN: produce_archive=False for a PDF
|
||||
WHEN: parse() is called
|
||||
THEN: the returned text is exactly the locally-extracted text,
|
||||
not anything from the (unused) Azure mock
|
||||
"""
|
||||
mocker.patch(
|
||||
"paperless.parsers.remote.extract_pdf_text",
|
||||
return_value="Local digital text.",
|
||||
)
|
||||
|
||||
remote_parser.parse(
|
||||
simple_digital_pdf_file,
|
||||
"application/pdf",
|
||||
produce_archive=False,
|
||||
)
|
||||
|
||||
assert remote_parser.get_text() == "Local digital text."
|
||||
|
||||
def test_pdf_no_archive_requested_closes_no_client(
|
||||
self,
|
||||
remote_parser: RemoteDocumentParser,
|
||||
simple_digital_pdf_file: Path,
|
||||
azure_client: Mock,
|
||||
) -> None:
|
||||
remote_parser.parse(
|
||||
simple_digital_pdf_file,
|
||||
"application/pdf",
|
||||
produce_archive=False,
|
||||
)
|
||||
|
||||
azure_client.close.assert_not_called()
|
||||
|
||||
def test_non_pdf_still_calls_azure_when_no_archive_requested(
|
||||
self,
|
||||
remote_parser: RemoteDocumentParser,
|
||||
simple_digital_pdf_file: Path,
|
||||
azure_client: Mock,
|
||||
) -> None:
|
||||
"""
|
||||
Images have no local-text fallback, so produce_archive=False does
|
||||
not skip the remote engine for non-PDF MIME types.
|
||||
"""
|
||||
remote_parser.parse(
|
||||
simple_digital_pdf_file,
|
||||
"image/png",
|
||||
produce_archive=False,
|
||||
)
|
||||
|
||||
azure_client.begin_analyze_document.assert_called_once()
|
||||
assert remote_parser.get_text() == _DEFAULT_TEXT
|
||||
|
||||
@pytest.mark.usefixtures("no_engine_settings")
|
||||
def test_unconfigured_engine_takes_precedence_over_skip(
|
||||
self,
|
||||
remote_parser: RemoteDocumentParser,
|
||||
simple_digital_pdf_file: Path,
|
||||
) -> None:
|
||||
"""An unconfigured engine still short-circuits before the
|
||||
produce_archive check, returning empty text as before.
|
||||
"""
|
||||
remote_parser.parse(
|
||||
simple_digital_pdf_file,
|
||||
"application/pdf",
|
||||
produce_archive=False,
|
||||
)
|
||||
|
||||
assert remote_parser.get_text() == ""
|
||||
assert remote_parser.get_archive_path() is None
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# parse() — Azure failure path
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@@ -21,7 +21,7 @@ from documents.parsers import run_convert
|
||||
from paperless.models import ModeChoices
|
||||
from paperless.parsers import ParserProtocol
|
||||
from paperless.parsers.tesseract import RasterisedDocumentParser
|
||||
from paperless.parsers.utils import post_process_text
|
||||
from paperless.parsers.tesseract import post_process_text
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from pathlib import Path
|
||||
|
||||
@@ -0,0 +1,229 @@
|
||||
"""Thin gateways over the plain relational side tables that sit alongside the
|
||||
vec0 table. Each method takes the sqlite3.Connection to operate on
|
||||
explicitly, rather than owning one -- the store swaps connections during
|
||||
compact()/migration, and migrations always work across two connections
|
||||
(src_conn, dst_conn) at once.
|
||||
|
||||
PRECONDITION: Callers must set conn.row_factory = sqlite3.Row before passing a
|
||||
connection to any of these gateways' read methods. The read methods across all
|
||||
three classes (DocumentChunksTable.chunk_ids_for_document, IndexMetaTable._get,
|
||||
DocumentMetaTable.all_modified_times, DocumentMetaTable.copy_all) use
|
||||
row["column_name"] dictionary-style indexing, which requires sqlite3.Row as the
|
||||
row factory -- without it, sqlite3.Row is not set, rows are returned as plain
|
||||
tuples, and tuple indices must be integers, raising TypeError.
|
||||
"""
|
||||
|
||||
import sqlite3
|
||||
from collections.abc import Iterable
|
||||
from typing import NamedTuple
|
||||
|
||||
|
||||
class ChunkRow(NamedTuple):
|
||||
chunk_id: str
|
||||
document_id: int
|
||||
|
||||
|
||||
class DocumentMetaRow(NamedTuple):
|
||||
document_id: int
|
||||
modified: str
|
||||
|
||||
|
||||
class DocumentChunksTable:
|
||||
"""chunk_id -> document_id, indexed by document_id. Gives O(1)
|
||||
per-document chunk lookup that vec0's own document_id metadata column
|
||||
cannot (see PaperlessSqliteVecVectorStore._delete_chunks_by_document_id).
|
||||
"""
|
||||
|
||||
@staticmethod
|
||||
def create(conn: sqlite3.Connection) -> None:
|
||||
conn.execute(
|
||||
"CREATE TABLE IF NOT EXISTS document_chunks "
|
||||
"(chunk_id TEXT PRIMARY KEY, document_id INTEGER NOT NULL)",
|
||||
)
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_document_chunks_document_id "
|
||||
"ON document_chunks (document_id)",
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def insert_many(conn: sqlite3.Connection, rows: Iterable[ChunkRow]) -> None:
|
||||
"""rows must already be batch-bounded by the caller (e.g. vec0's own
|
||||
fetchmany() loop) -- this never reads, so it can't itself introduce
|
||||
an unbounded scan, but a whole-table iterable defeats the point."""
|
||||
conn.executemany(
|
||||
"INSERT INTO document_chunks (chunk_id, document_id) VALUES (?, ?)",
|
||||
rows,
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def chunk_ids_for_document(
|
||||
conn: sqlite3.Connection,
|
||||
document_id: int,
|
||||
) -> list[str]:
|
||||
return [
|
||||
row["chunk_id"]
|
||||
for row in conn.execute(
|
||||
"SELECT chunk_id FROM document_chunks WHERE document_id = ?",
|
||||
(document_id,),
|
||||
).fetchall()
|
||||
]
|
||||
|
||||
@staticmethod
|
||||
def delete_for_document(conn: sqlite3.Connection, document_id: int) -> None:
|
||||
conn.execute(
|
||||
"DELETE FROM document_chunks WHERE document_id = ?",
|
||||
(document_id,),
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def delete_all(conn: sqlite3.Connection) -> None:
|
||||
conn.execute("DELETE FROM document_chunks")
|
||||
|
||||
@staticmethod
|
||||
def count(conn: sqlite3.Connection) -> int:
|
||||
"""Cheap stand-in for vec0's own row count -- see compact()."""
|
||||
return conn.execute("SELECT count(*) FROM document_chunks").fetchone()[0]
|
||||
|
||||
|
||||
class DocumentMetaTable:
|
||||
"""document_id -> modified, one row per document. Lives outside vec0
|
||||
because vec0 only inlines TEXT metadata up to 12 bytes and `modified`
|
||||
(an ISO timestamp) is always longer.
|
||||
"""
|
||||
|
||||
@staticmethod
|
||||
def create(conn: sqlite3.Connection) -> None:
|
||||
conn.execute(
|
||||
"CREATE TABLE IF NOT EXISTS document_meta "
|
||||
"(document_id INTEGER PRIMARY KEY, modified TEXT NOT NULL)",
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def upsert_many(
|
||||
conn: sqlite3.Connection,
|
||||
rows: Iterable[DocumentMetaRow],
|
||||
) -> None:
|
||||
conn.executemany(
|
||||
"INSERT INTO document_meta (document_id, modified) VALUES (?, ?) "
|
||||
"ON CONFLICT(document_id) DO UPDATE SET modified = excluded.modified",
|
||||
rows,
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def delete_for_document(conn: sqlite3.Connection, document_id: int) -> None:
|
||||
conn.execute(
|
||||
"DELETE FROM document_meta WHERE document_id = ?",
|
||||
(document_id,),
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def delete_all(conn: sqlite3.Connection) -> None:
|
||||
conn.execute("DELETE FROM document_meta")
|
||||
|
||||
@staticmethod
|
||||
def copy_all(
|
||||
src_conn: sqlite3.Connection,
|
||||
dst_conn: sqlite3.Connection,
|
||||
batch_size: int,
|
||||
) -> None:
|
||||
"""Stream document_meta from src_conn into dst_conn in bounded
|
||||
batches. The *only* sanctioned way to move this table across
|
||||
connections (compact()/migrations) -- an unbounded fetchall here
|
||||
would defeat the same OOM-avoidance the vec0 row copy already relies
|
||||
on. batch_size has no default: forces the call site to think about
|
||||
it (pass COMPACT_BATCH_SIZE)."""
|
||||
cursor = src_conn.execute(
|
||||
"SELECT document_id, modified FROM document_meta",
|
||||
)
|
||||
while batch := cursor.fetchmany(batch_size):
|
||||
DocumentMetaTable.upsert_many(
|
||||
dst_conn,
|
||||
(DocumentMetaRow(r["document_id"], r["modified"]) for r in batch),
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def all_modified_times(conn: sqlite3.Connection) -> dict[str, str]:
|
||||
"""Full document_id -> modified map, for get_modified_times()'s
|
||||
public API only. One unbounded read by design (existing behavior).
|
||||
Never use this for cross-connection copying; see copy_all()."""
|
||||
return {
|
||||
str(row["document_id"]): str(row["modified"] or "")
|
||||
for row in conn.execute(
|
||||
"SELECT document_id, modified FROM document_meta",
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
class IndexMetaTable:
|
||||
"""Typed accessors over index_meta's key/value rows -- replaces
|
||||
PaperlessSqliteVecVectorStore._meta_get_on/_meta_set_on, which returned
|
||||
untyped str | None regardless of whether the key held an int (dim,
|
||||
schema_version, total_inserts) or a string (embed_model).
|
||||
"""
|
||||
|
||||
@staticmethod
|
||||
def create(conn: sqlite3.Connection) -> None:
|
||||
conn.execute(
|
||||
"CREATE TABLE IF NOT EXISTS index_meta (key TEXT PRIMARY KEY, value TEXT)",
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _get(conn: sqlite3.Connection, key: str) -> str | None:
|
||||
row = conn.execute(
|
||||
"SELECT value FROM index_meta WHERE key = ?",
|
||||
(key,),
|
||||
).fetchone()
|
||||
return row["value"] if row else None
|
||||
|
||||
@staticmethod
|
||||
def _set(conn: sqlite3.Connection, key: str, value: str) -> None:
|
||||
conn.execute(
|
||||
"INSERT INTO index_meta (key, value) VALUES (?, ?) "
|
||||
"ON CONFLICT(key) DO UPDATE SET value = excluded.value",
|
||||
(key, value),
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def get_dim(conn: sqlite3.Connection) -> int | None:
|
||||
value = IndexMetaTable._get(conn, "dim")
|
||||
return int(value) if value is not None else None
|
||||
|
||||
@staticmethod
|
||||
def set_dim(conn: sqlite3.Connection, dim: int) -> None:
|
||||
IndexMetaTable._set(conn, "dim", str(dim))
|
||||
|
||||
@staticmethod
|
||||
def get_embed_model(conn: sqlite3.Connection) -> str | None:
|
||||
return IndexMetaTable._get(conn, "embed_model")
|
||||
|
||||
@staticmethod
|
||||
def set_embed_model(conn: sqlite3.Connection, name: str) -> None:
|
||||
IndexMetaTable._set(conn, "embed_model", name)
|
||||
|
||||
@staticmethod
|
||||
def get_schema_version(conn: sqlite3.Connection) -> int | None:
|
||||
value = IndexMetaTable._get(conn, "schema_version")
|
||||
return int(value) if value is not None else None
|
||||
|
||||
@staticmethod
|
||||
def set_schema_version(conn: sqlite3.Connection, version: int) -> None:
|
||||
IndexMetaTable._set(conn, "schema_version", str(version))
|
||||
|
||||
@staticmethod
|
||||
def get_total_inserts(conn: sqlite3.Connection) -> int:
|
||||
value = IndexMetaTable._get(conn, "total_inserts")
|
||||
return int(value) if value is not None else 0
|
||||
|
||||
@staticmethod
|
||||
def increment_total_inserts(conn: sqlite3.Connection, count: int) -> None:
|
||||
current = IndexMetaTable.get_total_inserts(conn)
|
||||
IndexMetaTable._set(conn, "total_inserts", str(current + count))
|
||||
|
||||
@staticmethod
|
||||
def reset_total_inserts(conn: sqlite3.Connection, count: int) -> None:
|
||||
"""Set total_inserts to an absolute value -- distinct from
|
||||
increment_total_inserts(): used by compact()'s rebuild and by
|
||||
m0001_v1_to_v2 after copying live rows into a fresh file, where
|
||||
total_inserts must become exactly the live row count, not add to
|
||||
whatever the source file's counter held."""
|
||||
IndexMetaTable._set(conn, "total_inserts", str(count))
|
||||
@@ -0,0 +1,306 @@
|
||||
import sqlite3
|
||||
from collections.abc import Generator
|
||||
|
||||
import pytest
|
||||
|
||||
from paperless_ai.tables import ChunkRow
|
||||
from paperless_ai.tables import DocumentChunksTable
|
||||
from paperless_ai.tables import DocumentMetaRow
|
||||
from paperless_ai.tables import DocumentMetaTable
|
||||
from paperless_ai.tables import IndexMetaTable
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def conn() -> Generator[sqlite3.Connection, None, None]:
|
||||
connection = sqlite3.connect(":memory:")
|
||||
connection.row_factory = sqlite3.Row
|
||||
try:
|
||||
yield connection
|
||||
finally:
|
||||
connection.close()
|
||||
|
||||
|
||||
class TestDocumentChunksTable:
|
||||
def test_create_is_idempotent(self, conn: sqlite3.Connection) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- A bare sqlite3 connection
|
||||
WHEN:
|
||||
- create() is called, a row is inserted, then create() is called again
|
||||
THEN:
|
||||
- No error is raised and the row survives uncorrupted
|
||||
"""
|
||||
DocumentChunksTable.create(conn)
|
||||
DocumentChunksTable.insert_many(conn, [ChunkRow("c1", 1)])
|
||||
DocumentChunksTable.create(conn)
|
||||
assert DocumentChunksTable.chunk_ids_for_document(conn, 1) == ["c1"]
|
||||
|
||||
def test_insert_many_then_lookup_by_document_id(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- An empty document_chunks table
|
||||
WHEN:
|
||||
- Two chunks for document 1 and one for document 2 are inserted
|
||||
THEN:
|
||||
- chunk_ids_for_document returns exactly the matching chunk ids
|
||||
"""
|
||||
DocumentChunksTable.create(conn)
|
||||
DocumentChunksTable.insert_many(
|
||||
conn,
|
||||
[ChunkRow("c1", 1), ChunkRow("c2", 1), ChunkRow("c3", 2)],
|
||||
)
|
||||
assert sorted(DocumentChunksTable.chunk_ids_for_document(conn, 1)) == [
|
||||
"c1",
|
||||
"c2",
|
||||
]
|
||||
assert DocumentChunksTable.chunk_ids_for_document(conn, 2) == ["c3"]
|
||||
assert DocumentChunksTable.chunk_ids_for_document(conn, 999) == []
|
||||
|
||||
def test_delete_for_document_removes_only_that_document(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- Chunks for two different documents
|
||||
WHEN:
|
||||
- delete_for_document() is called for one of them
|
||||
THEN:
|
||||
- Only that document's chunks are removed
|
||||
"""
|
||||
DocumentChunksTable.create(conn)
|
||||
DocumentChunksTable.insert_many(
|
||||
conn,
|
||||
[ChunkRow("c1", 1), ChunkRow("c2", 2)],
|
||||
)
|
||||
DocumentChunksTable.delete_for_document(conn, 1)
|
||||
assert DocumentChunksTable.chunk_ids_for_document(conn, 1) == []
|
||||
assert DocumentChunksTable.chunk_ids_for_document(conn, 2) == ["c2"]
|
||||
|
||||
def test_delete_all_clears_every_row(self, conn: sqlite3.Connection) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- Chunks for multiple documents
|
||||
WHEN:
|
||||
- delete_all() is called
|
||||
THEN:
|
||||
- count() returns 0
|
||||
"""
|
||||
DocumentChunksTable.create(conn)
|
||||
DocumentChunksTable.insert_many(
|
||||
conn,
|
||||
[ChunkRow("c1", 1), ChunkRow("c2", 2)],
|
||||
)
|
||||
DocumentChunksTable.delete_all(conn)
|
||||
assert DocumentChunksTable.count(conn) == 0
|
||||
|
||||
def test_count_reflects_live_rows(self, conn: sqlite3.Connection) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- An empty document_chunks table
|
||||
WHEN:
|
||||
- Rows are inserted then one document's rows are deleted
|
||||
THEN:
|
||||
- count() reflects the remaining row count
|
||||
"""
|
||||
DocumentChunksTable.create(conn)
|
||||
DocumentChunksTable.insert_many(
|
||||
conn,
|
||||
[ChunkRow("c1", 1), ChunkRow("c2", 1), ChunkRow("c3", 2)],
|
||||
)
|
||||
assert DocumentChunksTable.count(conn) == 3
|
||||
DocumentChunksTable.delete_for_document(conn, 1)
|
||||
assert DocumentChunksTable.count(conn) == 1
|
||||
|
||||
|
||||
class TestDocumentMetaTable:
|
||||
def test_upsert_many_then_all_modified_times(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- An empty document_meta table
|
||||
WHEN:
|
||||
- Two documents' modified timestamps are upserted
|
||||
THEN:
|
||||
- all_modified_times() returns both, keyed by str(document_id)
|
||||
"""
|
||||
DocumentMetaTable.create(conn)
|
||||
DocumentMetaTable.upsert_many(
|
||||
conn,
|
||||
[
|
||||
DocumentMetaRow(1, "2026-01-01T00:00:00"),
|
||||
DocumentMetaRow(2, "2026-02-02T00:00:00"),
|
||||
],
|
||||
)
|
||||
assert DocumentMetaTable.all_modified_times(conn) == {
|
||||
"1": "2026-01-01T00:00:00",
|
||||
"2": "2026-02-02T00:00:00",
|
||||
}
|
||||
|
||||
def test_upsert_many_overwrites_existing_value(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- A document_meta row for document 1
|
||||
WHEN:
|
||||
- upsert_many() is called again with a new modified value for
|
||||
the same document_id
|
||||
THEN:
|
||||
- The stored value is replaced, not duplicated
|
||||
"""
|
||||
DocumentMetaTable.create(conn)
|
||||
DocumentMetaTable.upsert_many(conn, [DocumentMetaRow(1, "old")])
|
||||
DocumentMetaTable.upsert_many(conn, [DocumentMetaRow(1, "new")])
|
||||
assert DocumentMetaTable.all_modified_times(conn) == {"1": "new"}
|
||||
|
||||
def test_delete_for_document_removes_only_that_row(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- document_meta rows for two documents
|
||||
WHEN:
|
||||
- delete_for_document() is called for one of them
|
||||
THEN:
|
||||
- Only that document's row is removed
|
||||
"""
|
||||
DocumentMetaTable.create(conn)
|
||||
DocumentMetaTable.upsert_many(
|
||||
conn,
|
||||
[DocumentMetaRow(1, "a"), DocumentMetaRow(2, "b")],
|
||||
)
|
||||
DocumentMetaTable.delete_for_document(conn, 1)
|
||||
assert DocumentMetaTable.all_modified_times(conn) == {"2": "b"}
|
||||
|
||||
def test_delete_all_clears_every_row(self, conn: sqlite3.Connection) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- document_meta rows for multiple documents
|
||||
WHEN:
|
||||
- delete_all() is called
|
||||
THEN:
|
||||
- all_modified_times() returns an empty dict
|
||||
"""
|
||||
DocumentMetaTable.create(conn)
|
||||
DocumentMetaTable.upsert_many(
|
||||
conn,
|
||||
[DocumentMetaRow(1, "a"), DocumentMetaRow(2, "b")],
|
||||
)
|
||||
DocumentMetaTable.delete_all(conn)
|
||||
assert DocumentMetaTable.all_modified_times(conn) == {}
|
||||
|
||||
def test_copy_all_streams_every_row_to_destination(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- A source connection with document_meta rows for 5 documents
|
||||
- A separate, empty destination connection
|
||||
WHEN:
|
||||
- copy_all() is called with a batch size smaller than the row
|
||||
count, forcing multiple fetchmany() cycles
|
||||
THEN:
|
||||
- Every row is present on the destination connection
|
||||
"""
|
||||
DocumentMetaTable.create(conn)
|
||||
DocumentMetaTable.upsert_many(
|
||||
conn,
|
||||
[DocumentMetaRow(i, f"modified-{i}") for i in range(5)],
|
||||
)
|
||||
dst_conn = sqlite3.connect(":memory:")
|
||||
dst_conn.row_factory = sqlite3.Row
|
||||
try:
|
||||
DocumentMetaTable.create(dst_conn)
|
||||
DocumentMetaTable.copy_all(conn, dst_conn, batch_size=2)
|
||||
assert DocumentMetaTable.all_modified_times(dst_conn) == {
|
||||
str(i): f"modified-{i}" for i in range(5)
|
||||
}
|
||||
finally:
|
||||
dst_conn.close()
|
||||
|
||||
|
||||
class TestIndexMetaTable:
|
||||
@pytest.mark.parametrize(
|
||||
("setter_name", "getter_name", "value"),
|
||||
[
|
||||
("set_dim", "get_dim", 384),
|
||||
("set_embed_model", "get_embed_model", "model-a"),
|
||||
("set_schema_version", "get_schema_version", 2),
|
||||
],
|
||||
)
|
||||
def test_typed_accessor_roundtrip(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
setter_name: str,
|
||||
getter_name: str,
|
||||
value: int | str,
|
||||
) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- An empty index_meta table
|
||||
WHEN:
|
||||
- A typed accessor's setter is called then the getter is read back
|
||||
THEN:
|
||||
- The same value is returned, correctly typed (int or str)
|
||||
"""
|
||||
IndexMetaTable.create(conn)
|
||||
getter = getattr(IndexMetaTable, getter_name)
|
||||
setter = getattr(IndexMetaTable, setter_name)
|
||||
assert getter(conn) is None
|
||||
setter(conn, value)
|
||||
assert getter(conn) == value
|
||||
|
||||
def test_total_inserts_starts_at_zero(self, conn: sqlite3.Connection) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- An empty index_meta table
|
||||
WHEN:
|
||||
- get_total_inserts() is read before anything is set
|
||||
THEN:
|
||||
- 0 is returned
|
||||
"""
|
||||
IndexMetaTable.create(conn)
|
||||
assert IndexMetaTable.get_total_inserts(conn) == 0
|
||||
|
||||
def test_increment_total_inserts_accumulates(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- An empty index_meta table
|
||||
WHEN:
|
||||
- increment_total_inserts() is called twice
|
||||
THEN:
|
||||
- get_total_inserts() returns the running sum
|
||||
"""
|
||||
IndexMetaTable.create(conn)
|
||||
IndexMetaTable.increment_total_inserts(conn, 5)
|
||||
IndexMetaTable.increment_total_inserts(conn, 3)
|
||||
assert IndexMetaTable.get_total_inserts(conn) == 8
|
||||
|
||||
def test_reset_total_inserts_sets_absolute_value(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""
|
||||
GIVEN:
|
||||
- A total_inserts counter already at a high value
|
||||
WHEN:
|
||||
- reset_total_inserts() is called with a lower value
|
||||
THEN:
|
||||
- get_total_inserts() returns exactly that value, not a sum
|
||||
"""
|
||||
IndexMetaTable.create(conn)
|
||||
IndexMetaTable.increment_total_inserts(conn, 100)
|
||||
IndexMetaTable.reset_total_inserts(conn, 7)
|
||||
assert IndexMetaTable.get_total_inserts(conn) == 7
|
||||
Reference in New Issue
Block a user