From f7b61f7da9e7a89743f538155ecf830e389f19bb Mon Sep 17 00:00:00 2001 From: oskar Date: Mon, 13 Jul 2026 20:00:00 +0200 Subject: [PATCH] =?UTF-8?q?feat(documents-ingest):=20PDF=20attachment=20ex?= =?UTF-8?q?tractor=20=E2=80=94=20mail=20archive=20->=20Paperless=20consume?= =?UTF-8?q?/?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Module 5 phase 1 (docs/kb/modules/05-documents-ingest.md): pulls a sample of PDF attachments (>50KB, last year, LIMIT 150) out of the Gmail .eml archive and drops them into Paperless consume/ for OCR, so the RAG layer has real documents to work with before the full Paperless/Nextcloud envelope adapter is built. sha256-first attachment matching (manifest filenames can carry raw RFC 2047 encoded-word artifacts that don't byte-match what email.policy.default decodes today — confirmed against live data, ~10% of candidates were affected). Idempotent via a sha256-keyed JSON registry; dry-run by default, --apply to write. Verified end-to-end on PIHA: dry-run + --apply both run against live kb-postgres/archive, 185/222 candidate PDFs written to consume/ (37 in-run duplicates correctly deduped), Paperless picked them up and started OCR immediately. --- jobs/documents-ingest/README.md | 183 ++++++++ jobs/documents-ingest/pyproject.toml | 22 + .../src/documents_ingest/__init__.py | 0 .../src/documents_ingest/extractor.py | 443 ++++++++++++++++++ jobs/documents-ingest/tests/test_extractor.py | 417 +++++++++++++++++ 5 files changed, 1065 insertions(+) create mode 100644 jobs/documents-ingest/README.md create mode 100644 jobs/documents-ingest/pyproject.toml create mode 100644 jobs/documents-ingest/src/documents_ingest/__init__.py create mode 100644 jobs/documents-ingest/src/documents_ingest/extractor.py create mode 100644 jobs/documents-ingest/tests/test_extractor.py diff --git a/jobs/documents-ingest/README.md b/jobs/documents-ingest/README.md new file mode 100644 index 0000000..5acf85c --- /dev/null +++ b/jobs/documents-ingest/README.md @@ -0,0 +1,183 @@ +# documents-ingest + +One-shot job CLI, Phase 1 of module 5 (`docs/kb/modules/05-documents-ingest.md`, +"Domkniecie dlugu z maili"). Extracts a **sample** of PDF attachments from the +Gmail `.eml` archive (already indexed in the `envelope` table of kb-postgres) +and drops them into Paperless' `consume/` directory so Paperless does the OCR +and correspondent-detection. This job does **not** write to the `envelope` +table — the Paperless/Nextcloud envelope adapter is a later phase of module 5. + +## Why a sample, not a bulk import + +The Gmail import left ~70k attachments referenced in `envelope.entities` +manifests (bytes live inside the archived `.eml` files, never extracted). +Dumping all of them into Paperless at once would swamp the OCR worker and the +RAG layer isn't built yet to make use of that volume. This job pulls a small, +recent, size-filtered sample (default: 150 envelopes, PDFs >50KB, from the +last year) as a testbed — mass import is a deliberate later decision. + +## Where it runs + +**Locally on PIHA**, as a plain CLI (not a container). It needs simultaneous +filesystem access to three things that all live on PIHA: + +- the mail archive (`/home/oskar/kb/mail/archive`) +- the Paperless `consume/` directory (`/opt/homelab/data/paperless/consume`) +- kb-postgres (`localhost:5433` from PIHA; reachable from elsewhere over + Tailscale, but the archive and consume dir are not — those are local paths) + +Install (from repo root, on PIHA): + +```bash +pip install -e jobs/documents-ingest/ +``` + +(Or reuse the venv already set up for `gmail-bulk-import`, e.g. +`/home/oskar/kb/venv/` — it already has `asyncpg` + `structlog`.) + +## Usage + +```bash +# Dry run (default) — preview only, no writes: +documents-ingest --dsn postgresql://kb:@localhost:5433/kb + +# Or via env var instead of --dsn: +export KB_DSN=postgresql://kb:@localhost:5433/kb +documents-ingest + +# Real run — write files into consume/ and update the registry: +documents-ingest --apply + +# Smaller/larger sample, different window/threshold: +documents-ingest --limit 50 --since-days 180 --min-size 100000 +``` + +Dry-run is the default and does not require `--consume-dir` to exist yet; +`--apply` does (Paperless must already be deployed with its consume dir in +place). See `documents-ingest --help` for all flags. + +## Candidate selection + +```sql +SELECT id, raw_ref, ts, entities FROM envelope +WHERE source = 'gmail' + AND ts > now() - interval '1 year' + AND EXISTS ( + SELECT 1 FROM jsonb_array_elements(entities) AS att + WHERE att->>'content_type' = 'application/pdf' + AND (att->>'size')::numeric > 50000 + ) +ORDER BY ts DESC +LIMIT 150 +``` + +For each matching envelope, every attachment manifest entry that passes the +filter is a separate candidate (one envelope can yield several PDFs). + +## Matching an attachment inside the .eml + +The manifest (`entities[]`) only has metadata — the attachment bytes live +inside the `.eml` (MIME multipart), so each candidate is resolved against the +freshly parsed message: + +1. Parse the `.eml` with `email.policy.default` and collect every + `application/pdf` MIME part (filename + decoded payload). +2. **sha256 is the proof of identity**, not the filename. The manifest was + built by a different parser at import time (`gmail-bulk-import`, using + `mailbox` + compat32 policy) and can still hold the raw RFC 2047 + encoded-word form of a filename (e.g. `=?UTF-8?b?...?=`, sometimes with + header-folding whitespace baked in), while `email.policy.default` decodes + it to real Unicode today. Comparing those byte-for-byte skipped ~10% of + otherwise-good attachments in testing — see `TestFindPdfParts` / + `TestProcessCandidate` in the test suite for the regression case. So: + match by sha256 across all PDF parts in the message; if none match, use a + filename match only to tell "found the named part but its bytes changed" + (`sha_mismatch`, reported and skipped) apart from "not present at all" + (`parse_error`, skipped). +3. The consume/ filename is built from the *decoded* filename (from the MIME + part), not the possibly-garbled manifest one. + +Mismatches and parse errors are never guessed past — they're logged and +skipped. + +## consume/ filenames + +`_.pdf`, date = envelope `ts`. On collision +(same date + sanitized name already used in this run or already present in +`consume/`), an 8-hex sha256 prefix is appended: +`__.pdf`. + +Files are written with a best-effort `chown` to uid:gid `1000:1000` (the +Paperless container's `USERMAP_UID/GID`, see `services/paperless/README.md`) +so Paperless can read them. If the chown fails (e.g. the job isn't running as +root/uid 1000), a warning is logged but the run continues — the write itself +already succeeded; fix ownership/perms on `consume/` separately if needed. +PIHA's uid/gid convention across the fleet is tracked as its own tech-debt +item (see `docs/backlog/`), not solved here. + +## Idempotency — registry + +A JSON file at `/opt/homelab/data/documents-ingest/registry.json` (default, +override with `--registry`), keyed by attachment sha256: + +```json +{ + "": { + "envelope_id": "...", + "filename": "...", + "consume_name": "2026-06-09_invoice.pdf", + "size": 123456, + "ingested_at": "2026-07-13T19:35:16+00:00" + } +} +``` + +**Why a JSON file and not a kb-postgres table:** this is a one-shot sampling +tool for a bootstrapping phase, not a long-running service — a new table +would formalize infrastructure for something temporary. A flat file needs no +migration, is trivial to inspect (`jq`) or reset, and sits under +`/opt/homelab/data/` alongside other node-local state per the repo's runtime +path convention. If/when module 5's real Paperless/Nextcloud adapter phase +starts writing `envelope` rows for `source=paperless`, that's the natural +point to fold this into a proper DB-backed ingest log — re-litigate then, not +now. + +Re-running the job only ever *adds* to the registry (on `--apply`); it's +never consulted or mutated in dry-run mode beyond being read for the preview. + +## Dry-run output + +Logs one line per skip (`skip.duplicate` / `skip.sha_mismatch` / +`skip.parse_error`, with reason), a `summary` line with full counts +(`envelopes_scanned`, `pdf_candidates`, `extracted`, `skipped_duplicate`, +`skipped_sha_mismatch`, `skipped_parse_error`, `errors`), and up to 20 example +`(target_name, size, envelope_id)` rows so you can sanity-check filenames +before running `--apply`. + +## Verifying the result in Paperless + +After `--apply`: + +1. Paperless' consumer picks files up from `consume/` automatically (polling + or inotify, per its own config) — no action needed on this job's side. +2. Watch progress: Paperless UI → Documents (new items appear as OCR + finishes), or `docker logs -f paperless` on PIHA for consumer/OCR activity. +3. Cross-check count: number of new documents in Paperless should equal + `stats["extracted"]` from the `--apply` run's summary line. +4. Confirm idempotency: re-running `--apply` immediately after should report + `extracted: 0` and `skipped_duplicate` equal to the previous run's + `extracted` count — nothing new lands in `consume/`. + +## Tests + +```bash +pip install -e jobs/documents-ingest/ +cd jobs/documents-ingest && pytest +``` + +Pure unit tests, no DB or filesystem outside `tmp_path` required — `run()` is +tested by monkeypatching `asyncpg.connect` with an in-memory fake connection. +Covers: filename sanitization, consume-name collision handling, manifest +filtering, MIME PDF-part extraction (including the RFC 2047 decoding +mismatch), sha256 match/mismatch, duplicate detection, dry-run vs `--apply` +behavior, and multi-attachment envelopes. diff --git a/jobs/documents-ingest/pyproject.toml b/jobs/documents-ingest/pyproject.toml new file mode 100644 index 0000000..b496292 --- /dev/null +++ b/jobs/documents-ingest/pyproject.toml @@ -0,0 +1,22 @@ +[build-system] +requires = ["setuptools>=68"] +build-backend = "setuptools.build_meta" + +[project] +name = "documents-ingest" +version = "0.1.0" +requires-python = ">=3.11" +dependencies = [ + "asyncpg>=0.29", + "structlog>=24.1", +] + +[project.scripts] +documents-ingest = "documents_ingest.extractor:main" + +[tool.setuptools.packages.find] +where = ["src"] + +[tool.pytest.ini_options] +asyncio_mode = "auto" +testpaths = ["tests"] diff --git a/jobs/documents-ingest/src/documents_ingest/__init__.py b/jobs/documents-ingest/src/documents_ingest/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/jobs/documents-ingest/src/documents_ingest/extractor.py b/jobs/documents-ingest/src/documents_ingest/extractor.py new file mode 100644 index 0000000..c1abbbf --- /dev/null +++ b/jobs/documents-ingest/src/documents_ingest/extractor.py @@ -0,0 +1,443 @@ +"""PDF attachment extractor — kb-postgres mail archive -> Paperless consume/. + +Module 5 ("documents-ingest"), Phase 1 (docs/kb/modules/05-documents-ingest.md, +section "Domkniecie dlugu z maili"): pull a sample of PDF attachments out of the +Gmail .eml archive and drop them into Paperless' consume/ dir so Paperless does +the OCR + correspondent-detection. This is NOT the Paperless/Nextcloud envelope +adapter (that is a later phase of module 5) — this job only produces files for +Paperless to ingest; it does not write to the `envelope` table itself. + +Runs on PIHA (the only host with the mail archive, the Paperless consume dir, +and kb-postgres all reachable without extra hops): + +Install (from repo root): + pip install -e jobs/documents-ingest/ + +Usage: + # Dry run (default) — preview what would be extracted, no writes: + documents-ingest --dsn postgresql://kb:@localhost:5433/kb + + # Real run — write files into consume/ and update the registry: + documents-ingest --dsn postgresql://kb:@localhost:5433/kb --apply + + # Smaller/larger sample, different window: + documents-ingest --dsn ... --limit 50 --since-days 180 --min-size 100000 + +DSN can also come from the KB_DSN env var instead of --dsn. +""" +from __future__ import annotations + +import argparse +import asyncio +import email +import email.policy +import hashlib +import json +import os +import re +import sys +from dataclasses import dataclass +from datetime import datetime, timedelta, timezone +from pathlib import Path +from typing import Optional + +import asyncpg +import structlog + +_log = structlog.get_logger(__name__) + +DEFAULT_ARCHIVE_ROOT = Path("/home/oskar/kb/mail/archive") +DEFAULT_CONSUME_DIR = Path("/opt/homelab/data/paperless/consume") +DEFAULT_REGISTRY = Path("/opt/homelab/data/documents-ingest/registry.json") +DEFAULT_LIMIT = 150 +DEFAULT_MIN_SIZE = 50_000 +DEFAULT_SINCE_DAYS = 365 + +CONSUME_UID = 1000 +CONSUME_GID = 1000 + +_UNSAFE_NAME = re.compile(r"[^A-Za-z0-9._-]+") +_MAX_NAME_LEN = 150 + + +@dataclass +class Candidate: + """One PDF attachment manifest entry that passed the size/type filter.""" + + envelope_id: str + ts: datetime + filename: str + size: int + sha256: str + + +def sanitize_filename(filename: str, fallback: str = "attachment.pdf") -> str: + """Reduce a manifest filename to a safe, bounded basename ending in .pdf.""" + name = Path(filename or "").name # drop any path components + name = _UNSAFE_NAME.sub("_", name).strip("._") + if not name: + name = fallback + if not name.lower().endswith(".pdf"): + name = f"{name}.pdf" + if len(name) > _MAX_NAME_LEN: + name = f"{name[: _MAX_NAME_LEN - 4]}.pdf" + return name + + +def build_consume_name(ts: datetime, filename: str, sha256: str, used: set[str]) -> str: + """Build a unique, readable consume/ filename: _[_].pdf. + + `used` tracks names already claimed in this run (and pre-seeded with names + already present in consume/) so repeated filenames don't collide. + """ + safe = sanitize_filename(filename) + stem = safe[: -len(".pdf")] + base = f"{ts:%Y-%m-%d}_{stem}.pdf" + if base not in used: + used.add(base) + return base + disambiguated = f"{ts:%Y-%m-%d}_{stem}_{sha256[:8]}.pdf" + used.add(disambiguated) + return disambiguated + + +def pdf_candidates_from_entities( + envelope_id: str, ts: datetime, entities: list, min_size: int +) -> list[Candidate]: + """Filter an envelope's attachment manifest down to qualifying PDF entries.""" + candidates = [] + for entity in entities: + if entity.get("type") != "attachment": + continue + if entity.get("content_type") != "application/pdf": + continue + size = entity.get("size") or 0 + if size <= min_size: + continue + sha256 = entity.get("sha256") + if not sha256: + continue + candidates.append( + Candidate( + envelope_id=envelope_id, + ts=ts, + filename=entity.get("filename") or "", + size=size, + sha256=sha256, + ) + ) + return candidates + + +def find_pdf_parts(raw: bytes) -> list[tuple[Optional[str], bytes]]: + """Return (filename, payload) for every application/pdf MIME part in the message. + + filename is whatever `email.policy.default` decodes it to (RFC 2047 aware) — + it will NOT byte-for-byte equal a raw manifest filename that still carries + encoded-word/folding artifacts (the manifest was built by a different parser + at import time). sha256 is the authoritative match; see process_candidate(). + """ + msg = email.message_from_bytes(raw, policy=email.policy.default) + parts = [] + for part in msg.walk(): + if part.get_content_maintype() == "multipart": + continue + if part.get_content_type() != "application/pdf": + continue + payload = part.get_payload(decode=True) + if payload is not None: + parts.append((part.get_filename(), payload)) + return parts + + +def load_registry(path: Path) -> dict[str, dict]: + """Load the sha256 -> ingested-record idempotency registry. Missing file = empty.""" + if not path.is_file(): + return {} + return json.loads(path.read_text()) + + +def save_registry(path: Path, registry: dict[str, dict]) -> None: + """Persist the registry atomically (write-tmp + rename).""" + path.parent.mkdir(parents=True, exist_ok=True) + tmp = path.with_suffix(path.suffix + ".tmp") + tmp.write_text(json.dumps(registry, indent=2, sort_keys=True)) + tmp.replace(path) + + +def _decode_jsonb(value: object) -> object: + """asyncpg may return jsonb as a str or an already-decoded object.""" + if value is None: + return None + if isinstance(value, str): + return json.loads(value) + return value + + +def _existing_consume_names(consume_dir: Path) -> set[str]: + try: + return {p.name for p in consume_dir.iterdir() if p.is_file()} + except OSError: + return set() + + +async def fetch_candidate_envelopes( + conn: asyncpg.Connection, since: datetime, min_size: int, limit: int +) -> list[asyncpg.Record]: + """Envelopes carrying at least one qualifying PDF attachment, newest first.""" + return await conn.fetch( + """ + SELECT id, raw_ref, ts, entities + FROM envelope + WHERE source = 'gmail' + AND ts > $1 + AND EXISTS ( + SELECT 1 FROM jsonb_array_elements(entities) AS att + WHERE att->>'content_type' = 'application/pdf' + AND (att->>'size')::numeric > $2 + ) + ORDER BY ts DESC + LIMIT $3 + """, + since, + min_size, + limit, + ) + + +def process_candidate( + raw: bytes, + candidate: Candidate, + registry: dict[str, dict], + used_names: set[str], + apply: bool, +) -> tuple[str, Optional[dict]]: + """Resolve one candidate against the .eml bytes and the registry. + + Returns (status, example) where status is one of: + "duplicate", "sha_mismatch", "parse_error", "extracted". + `example` is a dict describing the (would-be) consume/ file, or None. + Does not touch the filesystem — writing is the caller's job (see run()). + """ + if candidate.sha256 in registry: + _log.info("skip.duplicate", envelope_id=candidate.envelope_id, filename=candidate.filename) + return "duplicate", None + + try: + pdf_parts = find_pdf_parts(raw) + except Exception: + _log.warning( + "skip.parse_error", envelope_id=candidate.envelope_id, filename=candidate.filename, exc_info=True + ) + return "parse_error", None + + if not pdf_parts: + _log.warning( + "skip.parse_error", + envelope_id=candidate.envelope_id, + filename=candidate.filename, + reason="no_pdf_parts_in_mime", + ) + return "parse_error", None + + # sha256 is the proof of identity (manifest filenames may carry raw RFC 2047 + # encoded-word/folding artifacts from the original import parser, so they + # won't always byte-match what email.policy.default decodes today — the hash + # doesn't have that problem). Filename match is used only to distinguish + # "found the named part but its bytes changed" from "not present at all". + matched = next( + ( + (fname, payload) + for fname, payload in pdf_parts + if hashlib.sha256(payload).hexdigest() == candidate.sha256 + ), + None, + ) + if matched is None: + if any(fname == candidate.filename for fname, _ in pdf_parts): + _log.warning( + "skip.sha_mismatch", + envelope_id=candidate.envelope_id, + filename=candidate.filename, + expected_sha256=candidate.sha256, + ) + return "sha_mismatch", None + _log.warning( + "skip.parse_error", + envelope_id=candidate.envelope_id, + filename=candidate.filename, + reason="attachment_not_found_in_mime", + ) + return "parse_error", None + + # Prefer the decoded filename from the MIME part itself for the human-facing + # consume/ name — the manifest filename may still carry encoded-word artifacts. + matched_filename, matched_payload = matched + display_filename = matched_filename or candidate.filename + target_name = build_consume_name(candidate.ts, display_filename, candidate.sha256, used_names) + example = { + "envelope_id": candidate.envelope_id, + "filename": display_filename, + "target_name": target_name, + "size": candidate.size, + "sha256": candidate.sha256, + "payload": matched_payload if apply else None, + } + return "extracted", example + + +async def run( + dsn: str, + archive_root: Path = DEFAULT_ARCHIVE_ROOT, + consume_dir: Path = DEFAULT_CONSUME_DIR, + registry_path: Path = DEFAULT_REGISTRY, + limit: int = DEFAULT_LIMIT, + min_size: int = DEFAULT_MIN_SIZE, + since_days: int = DEFAULT_SINCE_DAYS, + apply: bool = False, +) -> tuple[dict[str, int], list[dict]]: + """Scan candidate envelopes and (if apply) extract PDFs into consume/. + + Returns (stats, examples). stats always reflects what was found/matched; + in dry-run mode "extracted" counts what *would* be written. + """ + conn = await asyncpg.connect(dsn) + try: + since = datetime.now(timezone.utc) - timedelta(days=since_days) + rows = await fetch_candidate_envelopes(conn, since, min_size, limit) + finally: + await conn.close() + + registry = load_registry(registry_path) + used_names = _existing_consume_names(consume_dir) + + stats = { + "envelopes_scanned": 0, + "pdf_candidates": 0, + "extracted": 0, + "skipped_duplicate": 0, + "skipped_sha_mismatch": 0, + "skipped_parse_error": 0, + "errors": 0, + } + examples: list[dict] = [] + + for row in rows: + stats["envelopes_scanned"] += 1 + entities = _decode_jsonb(row["entities"]) or [] + ts = row["ts"] + if ts.tzinfo is None: + ts = ts.replace(tzinfo=timezone.utc) + + candidates = pdf_candidates_from_entities(row["id"], ts, entities, min_size) + if not candidates: + continue + + eml_path = archive_root / row["raw_ref"] + try: + raw = eml_path.read_bytes() + except OSError: + _log.error("eml_read_failed", envelope_id=row["id"], path=str(eml_path)) + stats["errors"] += len(candidates) + continue + + for candidate in candidates: + stats["pdf_candidates"] += 1 + status, example = process_candidate(raw, candidate, registry, used_names, apply) + stats[ + { + "duplicate": "skipped_duplicate", + "sha_mismatch": "skipped_sha_mismatch", + "parse_error": "skipped_parse_error", + "extracted": "extracted", + }[status] + ] += 1 + + if status != "extracted" or example is None: + continue + examples.append({k: v for k, v in example.items() if k != "payload"}) + + if not apply: + continue + + target_path = consume_dir / example["target_name"] + target_path.write_bytes(example["payload"]) + try: + os.chown(target_path, CONSUME_UID, CONSUME_GID) + except PermissionError: + _log.warning("chown_failed", path=str(target_path)) + + registry[candidate.sha256] = { + "envelope_id": candidate.envelope_id, + "filename": example["filename"], + "consume_name": example["target_name"], + "size": candidate.size, + "ingested_at": datetime.now(timezone.utc).isoformat(), + } + _log.info("written", envelope_id=candidate.envelope_id, target=example["target_name"]) + + if apply: + save_registry(registry_path, registry) + + _log.info("run_complete", apply=apply, **stats) + return stats, examples + + +def main() -> None: + parser = argparse.ArgumentParser( + description="Extract PDF attachments from the mail archive into Paperless consume/ " + "(module 5, phase 1 — sample, not bulk)." + ) + parser.add_argument("--dsn", default=os.environ.get("KB_DSN"), + help="asyncpg DSN for kb-postgres (or set KB_DSN env var)") + parser.add_argument("--archive-root", type=Path, default=DEFAULT_ARCHIVE_ROOT, + help=f"Mail .eml archive root (default: {DEFAULT_ARCHIVE_ROOT})") + parser.add_argument("--consume-dir", type=Path, default=DEFAULT_CONSUME_DIR, + help=f"Paperless consume/ directory (default: {DEFAULT_CONSUME_DIR})") + parser.add_argument("--registry", type=Path, default=DEFAULT_REGISTRY, + help=f"Idempotency registry JSON (default: {DEFAULT_REGISTRY})") + parser.add_argument("--limit", type=int, default=DEFAULT_LIMIT, + help=f"Max envelopes to consider (default: {DEFAULT_LIMIT} — a sample, not bulk)") + parser.add_argument("--min-size", type=int, default=DEFAULT_MIN_SIZE, + help=f"Minimum PDF attachment size in bytes (default: {DEFAULT_MIN_SIZE})") + parser.add_argument("--since-days", type=int, default=DEFAULT_SINCE_DAYS, + help=f"Only envelopes newer than N days (default: {DEFAULT_SINCE_DAYS})") + parser.add_argument("--apply", action="store_true", + help="Actually write files and update the registry. Default is dry-run.") + args = parser.parse_args() + + if not args.dsn: + _log.error("missing_dsn", hint="pass --dsn or set KB_DSN") + sys.exit(1) + if not args.archive_root.is_dir(): + _log.error("archive_root_not_found", path=str(args.archive_root)) + sys.exit(1) + if args.apply and not args.consume_dir.is_dir(): + _log.error("consume_dir_not_found", path=str(args.consume_dir)) + sys.exit(1) + + stats, examples = asyncio.run( + run( + dsn=args.dsn, + archive_root=args.archive_root, + consume_dir=args.consume_dir, + registry_path=args.registry, + limit=args.limit, + min_size=args.min_size, + since_days=args.since_days, + apply=args.apply, + ) + ) + + mode = "APPLY" if args.apply else "DRY-RUN" + _log.info("summary", mode=mode, **stats) + for ex in examples[:20]: + _log.info("example", target_name=ex["target_name"], size=ex["size"], envelope_id=ex["envelope_id"]) + if len(examples) > 20: + _log.info("examples_truncated", shown=20, total=len(examples)) + + sys.exit(1 if stats["errors"] > 0 else 0) + + +if __name__ == "__main__": + main() diff --git a/jobs/documents-ingest/tests/test_extractor.py b/jobs/documents-ingest/tests/test_extractor.py new file mode 100644 index 0000000..4c75c73 --- /dev/null +++ b/jobs/documents-ingest/tests/test_extractor.py @@ -0,0 +1,417 @@ +"""Unit tests for the documents-ingest PDF extractor — no DB, no external services.""" +from __future__ import annotations + +import hashlib +import json +from datetime import datetime, timezone +from email import encoders +from email.mime.application import MIMEApplication +from email.mime.multipart import MIMEMultipart +from email.mime.text import MIMEText + +import pytest + +from documents_ingest.extractor import ( + Candidate, + build_consume_name, + find_pdf_parts, + load_registry, + pdf_candidates_from_entities, + process_candidate, + run, + sanitize_filename, + save_registry, +) + + +def _build_eml(parts: list[tuple[str, str, bytes]]) -> bytes: + """parts: list of (filename, content_type, bytes). Builds a multipart .eml.""" + msg = MIMEMultipart() + msg["From"] = "alice@example.com" + msg["To"] = "bob@example.com" + msg["Subject"] = "Invoice attached" + msg["Message-ID"] = "" + msg["Date"] = "Tue, 10 Jun 2025 12:00:00 +0000" + msg.attach(MIMEText("See attached", "plain")) + for filename, content_type, data in parts: + maintype, subtype = content_type.split("/", 1) + att = MIMEApplication(data, _subtype=subtype) if maintype == "application" else None + if att is None: + from email.mime.base import MIMEBase + att = MIMEBase(maintype, subtype) + att.set_payload(data) + encoders.encode_base64(att) + att.add_header("Content-Disposition", "attachment", filename=filename) + msg.attach(att) + return msg.as_bytes() + + +PDF_BYTES = b"%PDF-1.4 fake pdf content for testing" +PDF_SHA256 = hashlib.sha256(PDF_BYTES).hexdigest() + + +class TestSanitizeFilename: + def test_strips_path_components(self): + assert sanitize_filename("../../etc/passwd.pdf") == "passwd.pdf" + + def test_replaces_unsafe_chars(self): + assert sanitize_filename("faktura #123 (final).pdf") == "faktura_123_final_.pdf" + + def test_appends_pdf_extension_if_missing(self): + assert sanitize_filename("report").endswith(".pdf") + + def test_empty_name_falls_back(self): + assert sanitize_filename("") == "attachment.pdf" + + def test_only_unsafe_chars_falls_back(self): + assert sanitize_filename("///...") == "attachment.pdf" + + def test_truncates_long_names(self): + long_name = ("a" * 300) + ".pdf" + result = sanitize_filename(long_name) + assert len(result) <= 150 + assert result.endswith(".pdf") + + def test_preserves_reasonable_name(self): + assert sanitize_filename("faktura_2026-06.pdf") == "faktura_2026-06.pdf" + + +class TestBuildConsumeName: + def test_basic_name(self): + ts = datetime(2026, 6, 1, tzinfo=timezone.utc) + used: set[str] = set() + name = build_consume_name(ts, "invoice.pdf", "abc123", used) + assert name == "2026-06-01_invoice.pdf" + + def test_collision_gets_hash_suffix(self): + ts = datetime(2026, 6, 1, tzinfo=timezone.utc) + used: set[str] = set() + first = build_consume_name(ts, "invoice.pdf", "aaaa1111", used) + second = build_consume_name(ts, "invoice.pdf", "bbbb2222", used) + assert first != second + assert second == "2026-06-01_invoice_bbbb2222.pdf" + + def test_names_are_registered_in_used_set(self): + ts = datetime(2026, 6, 1, tzinfo=timezone.utc) + used: set[str] = set() + name = build_consume_name(ts, "invoice.pdf", "abc123", used) + assert name in used + + +class TestPdfCandidatesFromEntities: + ts = datetime(2026, 6, 1, tzinfo=timezone.utc) + + def test_filters_by_content_type_and_size(self): + entities = [ + {"type": "attachment", "content_type": "application/pdf", "size": 60000, + "filename": "big.pdf", "sha256": "x"}, + {"type": "attachment", "content_type": "application/pdf", "size": 1000, + "filename": "small.pdf", "sha256": "y"}, + {"type": "attachment", "content_type": "image/png", "size": 100000, + "filename": "photo.png", "sha256": "z"}, + ] + candidates = pdf_candidates_from_entities("env1", self.ts, entities, min_size=50000) + assert len(candidates) == 1 + assert candidates[0].filename == "big.pdf" + + def test_skips_entries_without_sha256(self): + entities = [ + {"type": "attachment", "content_type": "application/pdf", "size": 60000, + "filename": "big.pdf"}, + ] + assert pdf_candidates_from_entities("env1", self.ts, entities, min_size=50000) == [] + + def test_ignores_non_attachment_entries(self): + entities = [{"type": "inline", "content_type": "application/pdf", "size": 60000}] + assert pdf_candidates_from_entities("env1", self.ts, entities, min_size=50000) == [] + + +class TestFindPdfParts: + def test_finds_pdf_part(self): + raw = _build_eml([("report.pdf", "application/pdf", PDF_BYTES)]) + parts = find_pdf_parts(raw) + assert parts == [("report.pdf", PDF_BYTES)] + + def test_no_pdf_parts_returns_empty(self): + raw = _build_eml([("report.png", "image/png", b"not a pdf")]) + assert find_pdf_parts(raw) == [] + + def test_multiple_pdf_attachments_all_returned(self): + raw = _build_eml([ + ("a.pdf", "application/pdf", PDF_BYTES), + ("b.pdf", "application/pdf", b"other pdf"), + ("photo.png", "image/png", b"png bytes"), + ]) + parts = find_pdf_parts(raw) + assert ("a.pdf", PDF_BYTES) in parts + assert ("b.pdf", b"other pdf") in parts + assert len(parts) == 2 + + def test_decodes_rfc2047_encoded_filenames(self): + # Regression: the manifest (built by a different parser at import time) + # stores the RAW encoded-word filename; email.policy.default decodes it + # to real unicode here. process_candidate() must not require these to + # match byte-for-byte — sha256 is the authority (see TestProcessCandidate). + raw = _build_eml([ + ("=?UTF-8?b?ZmFrdHVyYV/FunJvZGxvLnBkZg==?=", "application/pdf", PDF_BYTES), + ]) + parts = find_pdf_parts(raw) + assert len(parts) == 1 + filename, payload = parts[0] + assert filename == "faktura_źrodlo.pdf" + assert payload == PDF_BYTES + + +class TestProcessCandidate: + ts = datetime(2026, 6, 1, tzinfo=timezone.utc) + + def _candidate(self, sha256=PDF_SHA256, filename="report.pdf"): + return Candidate(envelope_id="env1", ts=self.ts, filename=filename, + size=len(PDF_BYTES), sha256=sha256) + + def test_extracted_on_sha_match(self): + raw = _build_eml([("report.pdf", "application/pdf", PDF_BYTES)]) + status, example = process_candidate(raw, self._candidate(), {}, set(), apply=False) + assert status == "extracted" + assert example["target_name"] == "2026-06-01_report.pdf" + assert example["sha256"] == PDF_SHA256 + + def test_duplicate_when_sha_in_registry(self): + raw = _build_eml([("report.pdf", "application/pdf", PDF_BYTES)]) + registry = {PDF_SHA256: {"envelope_id": "env0"}} + status, example = process_candidate(raw, self._candidate(), registry, set(), apply=False) + assert status == "duplicate" + assert example is None + + def test_sha_mismatch_reported_and_skipped(self): + raw = _build_eml([("report.pdf", "application/pdf", PDF_BYTES)]) + candidate = self._candidate(sha256="0" * 64) + status, example = process_candidate(raw, candidate, {}, set(), apply=False) + assert status == "sha_mismatch" + assert example is None + + def test_parse_error_when_attachment_missing_from_mime(self): + raw = _build_eml([("other.pdf", "application/pdf", b"unrelated pdf bytes")]) + status, example = process_candidate(raw, self._candidate(), {}, set(), apply=False) + assert status == "parse_error" + assert example is None + + def test_sha_match_wins_over_filename_mismatch(self): + # Regression: manifest filename is the raw RFC 2047 encoded-word form, + # but the freshly re-parsed MIME part decodes to real unicode text. + # sha256 is the proof of identity — this must still extract. + raw = _build_eml([ + ("=?UTF-8?b?ZmFrdHVyYV/FunJvZGxvLnBkZg==?=", "application/pdf", PDF_BYTES), + ]) + candidate = self._candidate(filename="=?UTF-8?b?ZmFrdHVyYV/FunJvZGxvLnBkZg==?=") + status, example = process_candidate(raw, candidate, {}, set(), apply=False) + assert status == "extracted" + assert example["sha256"] == PDF_SHA256 + # Consume/ filename is built from the decoded name, not the raw encoded-word. + assert example["filename"] == "faktura_źrodlo.pdf" + assert "=?UTF-8?" not in example["target_name"] + + def test_apply_mode_includes_payload_for_writing(self): + raw = _build_eml([("report.pdf", "application/pdf", PDF_BYTES)]) + status, example = process_candidate(raw, self._candidate(), {}, set(), apply=True) + assert status == "extracted" + assert example["payload"] == PDF_BYTES + + def test_dry_run_mode_has_no_payload(self): + raw = _build_eml([("report.pdf", "application/pdf", PDF_BYTES)]) + status, example = process_candidate(raw, self._candidate(), {}, set(), apply=False) + assert example["payload"] is None + + +class TestRegistry: + def test_missing_file_returns_empty_dict(self, tmp_path): + assert load_registry(tmp_path / "nope.json") == {} + + def test_round_trip(self, tmp_path): + path = tmp_path / "sub" / "registry.json" + registry = {"abc123": {"envelope_id": "env1", "consume_name": "x.pdf"}} + save_registry(path, registry) + assert load_registry(path) == registry + + def test_save_creates_parent_dirs(self, tmp_path): + path = tmp_path / "a" / "b" / "registry.json" + save_registry(path, {}) + assert path.is_file() + + +class _FakeConn: + def __init__(self, rows): + self._rows = rows + + async def fetch(self, query, *params): + return self._rows + + async def close(self): + pass + + +def _row(envelope_id, ts, raw_ref, entities): + return {"id": envelope_id, "raw_ref": raw_ref, "ts": ts, "entities": json.dumps(entities)} + + +def _pdf_entity(filename="report.pdf", size=60000, sha256=PDF_SHA256): + return {"type": "attachment", "content_type": "application/pdf", "size": size, + "filename": filename, "sha256": sha256} + + +class TestRun: + ts = datetime(2026, 6, 1, tzinfo=timezone.utc) + + def _setup_archive(self, tmp_path, raw_ref, raw_bytes): + archive_root = tmp_path / "archive" + eml_path = archive_root / raw_ref + eml_path.parent.mkdir(parents=True, exist_ok=True) + eml_path.write_bytes(raw_bytes) + return archive_root + + def _patch_connect(self, monkeypatch, rows): + async def _fake_connect(dsn): + return _FakeConn(rows) + monkeypatch.setattr("documents_ingest.extractor.asyncpg.connect", _fake_connect) + + async def test_dry_run_does_not_write_files_or_registry(self, tmp_path, monkeypatch): + raw = _build_eml([("report.pdf", "application/pdf", PDF_BYTES)]) + archive_root = self._setup_archive(tmp_path, "gmail/2026/06/m1.eml", raw) + consume_dir = tmp_path / "consume" + consume_dir.mkdir() + registry_path = tmp_path / "registry.json" + rows = [_row("env1", self.ts, "gmail/2026/06/m1.eml", [_pdf_entity()])] + self._patch_connect(monkeypatch, rows) + + stats, examples = await run( + dsn="postgresql://fake", archive_root=archive_root, consume_dir=consume_dir, + registry_path=registry_path, apply=False, + ) + + assert stats["extracted"] == 1 + assert len(examples) == 1 + assert list(consume_dir.iterdir()) == [] + assert not registry_path.exists() + + async def test_apply_writes_file_and_registry(self, tmp_path, monkeypatch): + raw = _build_eml([("report.pdf", "application/pdf", PDF_BYTES)]) + archive_root = self._setup_archive(tmp_path, "gmail/2026/06/m1.eml", raw) + consume_dir = tmp_path / "consume" + consume_dir.mkdir() + registry_path = tmp_path / "registry.json" + rows = [_row("env1", self.ts, "gmail/2026/06/m1.eml", [_pdf_entity()])] + self._patch_connect(monkeypatch, rows) + + stats, examples = await run( + dsn="postgresql://fake", archive_root=archive_root, consume_dir=consume_dir, + registry_path=registry_path, apply=True, + ) + + assert stats["extracted"] == 1 + written = list(consume_dir.iterdir()) + assert len(written) == 1 + assert written[0].read_bytes() == PDF_BYTES + registry = load_registry(registry_path) + assert PDF_SHA256 in registry + + async def test_idempotent_second_apply_run_skips_duplicate(self, tmp_path, monkeypatch): + raw = _build_eml([("report.pdf", "application/pdf", PDF_BYTES)]) + archive_root = self._setup_archive(tmp_path, "gmail/2026/06/m1.eml", raw) + consume_dir = tmp_path / "consume" + consume_dir.mkdir() + registry_path = tmp_path / "registry.json" + rows = [_row("env1", self.ts, "gmail/2026/06/m1.eml", [_pdf_entity()])] + self._patch_connect(monkeypatch, rows) + + await run(dsn="postgresql://fake", archive_root=archive_root, consume_dir=consume_dir, + registry_path=registry_path, apply=True) + stats2, examples2 = await run( + dsn="postgresql://fake", archive_root=archive_root, consume_dir=consume_dir, + registry_path=registry_path, apply=True, + ) + + assert stats2["extracted"] == 0 + assert stats2["skipped_duplicate"] == 1 + assert len(list(consume_dir.iterdir())) == 1 + + async def test_multiple_candidates_in_one_envelope(self, tmp_path, monkeypatch): + pdf2 = b"%PDF-1.4 second document" + sha2 = hashlib.sha256(pdf2).hexdigest() + raw = _build_eml([ + ("a.pdf", "application/pdf", PDF_BYTES), + ("b.pdf", "application/pdf", pdf2), + ]) + archive_root = self._setup_archive(tmp_path, "gmail/2026/06/m1.eml", raw) + consume_dir = tmp_path / "consume" + consume_dir.mkdir() + registry_path = tmp_path / "registry.json" + rows = [_row("env1", self.ts, "gmail/2026/06/m1.eml", [ + _pdf_entity(filename="a.pdf", sha256=PDF_SHA256), + _pdf_entity(filename="b.pdf", sha256=sha2), + ])] + self._patch_connect(monkeypatch, rows) + + stats, examples = await run( + dsn="postgresql://fake", archive_root=archive_root, consume_dir=consume_dir, + registry_path=registry_path, apply=True, + ) + + assert stats["extracted"] == 2 + assert len(list(consume_dir.iterdir())) == 2 + + async def test_sha_mismatch_is_not_written(self, tmp_path, monkeypatch): + raw = _build_eml([("report.pdf", "application/pdf", PDF_BYTES)]) + archive_root = self._setup_archive(tmp_path, "gmail/2026/06/m1.eml", raw) + consume_dir = tmp_path / "consume" + consume_dir.mkdir() + registry_path = tmp_path / "registry.json" + rows = [_row("env1", self.ts, "gmail/2026/06/m1.eml", [ + _pdf_entity(sha256="0" * 64), + ])] + self._patch_connect(monkeypatch, rows) + + stats, examples = await run( + dsn="postgresql://fake", archive_root=archive_root, consume_dir=consume_dir, + registry_path=registry_path, apply=True, + ) + + assert stats["extracted"] == 0 + assert stats["skipped_sha_mismatch"] == 1 + assert list(consume_dir.iterdir()) == [] + + async def test_missing_eml_file_counts_as_error(self, tmp_path, monkeypatch): + archive_root = tmp_path / "archive" + archive_root.mkdir() + consume_dir = tmp_path / "consume" + consume_dir.mkdir() + registry_path = tmp_path / "registry.json" + rows = [_row("env1", self.ts, "gmail/2026/06/missing.eml", [_pdf_entity()])] + self._patch_connect(monkeypatch, rows) + + stats, examples = await run( + dsn="postgresql://fake", archive_root=archive_root, consume_dir=consume_dir, + registry_path=registry_path, apply=True, + ) + + assert stats["errors"] == 1 + assert stats["extracted"] == 0 + + async def test_envelope_without_qualifying_entities_is_scanned_but_skipped(self, tmp_path, monkeypatch): + archive_root = tmp_path / "archive" + archive_root.mkdir() + consume_dir = tmp_path / "consume" + consume_dir.mkdir() + registry_path = tmp_path / "registry.json" + rows = [_row("env1", self.ts, "gmail/2026/06/m1.eml", [ + {"type": "attachment", "content_type": "image/png", "size": 100, "filename": "x.png", "sha256": "z"}, + ])] + self._patch_connect(monkeypatch, rows) + + stats, examples = await run( + dsn="postgresql://fake", archive_root=archive_root, consume_dir=consume_dir, + registry_path=registry_path, apply=False, + ) + + assert stats["envelopes_scanned"] == 1 + assert stats["pdf_candidates"] == 0 + assert examples == []