diff --git a/jobs/gmail-bulk-import/Dockerfile b/jobs/gmail-bulk-import/Dockerfile new file mode 100644 index 0000000..a81cd62 --- /dev/null +++ b/jobs/gmail-bulk-import/Dockerfile @@ -0,0 +1,11 @@ +FROM python:3.11-slim + +WORKDIR /app + +COPY packages/kb-mail/ /packages/kb-mail/ +RUN pip install --no-cache-dir /packages/kb-mail/ + +COPY jobs/gmail-bulk-import/ /app/ +RUN pip install --no-cache-dir /app/ + +ENTRYPOINT ["gmail-bulk-import"] diff --git a/jobs/gmail-bulk-import/docker-compose.yml b/jobs/gmail-bulk-import/docker-compose.yml new file mode 100644 index 0000000..d90fc01 --- /dev/null +++ b/jobs/gmail-bulk-import/docker-compose.yml @@ -0,0 +1,13 @@ +services: + gmail-bulk-import: + build: + context: ../.. + dockerfile: jobs/gmail-bulk-import/Dockerfile + restart: "no" + volumes: + - ${MBOX_PATH:?set MBOX_PATH=/path/to/allmail.mbox}:/data/input.mbox:ro + - ${ARCHIVE_ROOT:?set ARCHIVE_ROOT=/path/to/archive}:/data/archive + command: + - --mbox=/data/input.mbox + - --archive=/data/archive + - --dsn=${KB_DSN:-} diff --git a/jobs/gmail-bulk-import/pyproject.toml b/jobs/gmail-bulk-import/pyproject.toml new file mode 100644 index 0000000..39c467e --- /dev/null +++ b/jobs/gmail-bulk-import/pyproject.toml @@ -0,0 +1,23 @@ +[build-system] +requires = ["setuptools>=68"] +build-backend = "setuptools.build_meta" + +[project] +name = "gmail-bulk-import" +version = "0.1.0" +requires-python = ">=3.11" +dependencies = [ + "asyncpg>=0.29", + "structlog>=24.1", + "kb-mail", +] + +[project.scripts] +gmail-bulk-import = "gmail_bulk_import.importer:main" + +[tool.setuptools.packages.find] +where = ["src"] + +[tool.pytest.ini_options] +asyncio_mode = "auto" +testpaths = ["tests"] diff --git a/jobs/gmail-bulk-import/src/gmail_bulk_import/__init__.py b/jobs/gmail-bulk-import/src/gmail_bulk_import/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/jobs/gmail-bulk-import/src/gmail_bulk_import/importer.py b/jobs/gmail-bulk-import/src/gmail_bulk_import/importer.py new file mode 100644 index 0000000..7afee38 --- /dev/null +++ b/jobs/gmail-bulk-import/src/gmail_bulk_import/importer.py @@ -0,0 +1,146 @@ +"""Gmail Takeout bulk importer — one-shot job. + +Reads a .mbox file (Google Takeout "All Mail" export) and imports each +message into the append-only .eml archive and optionally into kb-postgres. + +Idempotent: re-running against the same archive skips already-imported +messages (FileExistsError from archive → skipped; ON CONFLICT DO NOTHING in DB). +""" +from __future__ import annotations + +import argparse +import asyncio +import email.policy +import hashlib +import mailbox +import sys +from datetime import datetime, timezone +from email.utils import parsedate_to_datetime +from pathlib import Path +from typing import Optional + +import asyncpg +import structlog + +from kb_mail.archive import save_eml +from kb_mail.db import insert_envelope +from kb_mail.envelope import Envelope + +_log = structlog.get_logger(__name__) + + +def _message_id(msg: mailbox.mboxMessage) -> str: + """Return a stable envelope id from Message-ID header or SHA-256 content hash.""" + mid = msg.get("Message-ID", "").strip().strip("<>") + if mid: + return mid + return "sha256-" + hashlib.sha256(msg.as_bytes()).hexdigest()[:32] + + +def _parse_date(msg: mailbox.mboxMessage) -> datetime: + """Parse Date header to UTC-aware datetime. Falls back to epoch if unparseable.""" + date_str = msg.get("Date", "") + if date_str: + try: + dt = parsedate_to_datetime(date_str) + if dt.tzinfo is None: + dt = dt.replace(tzinfo=timezone.utc) + return dt.astimezone(timezone.utc) + except Exception: + _log.warning("unparseable_date", header=date_str[:120]) + return datetime(1970, 1, 1, tzinfo=timezone.utc) + + +async def _import_one( + msg: mailbox.mboxMessage, + archive_root: Path, + conn: Optional[asyncpg.Connection], +) -> str: + """Import a single message. Returns 'imported', 'skipped', or 'error'.""" + try: + envelope_id = _message_id(msg) + ts = _parse_date(msg) + raw = msg.as_bytes(policy=email.policy.compat32) + + try: + raw_ref = await save_eml(archive_root, envelope_id, "gmail", ts, raw) + except FileExistsError: + return "skipped" + + if conn is not None: + env = Envelope(id=envelope_id, source="gmail", ts=ts, raw_ref=raw_ref) + await insert_envelope(conn, env) + + return "imported" + except Exception: + _log.exception("message_failed", id=msg.get("Message-ID", "")) + return "error" + + +async def run_import( + mbox_path: Path, + archive_root: Path, + dsn: Optional[str] = None, + dry_run: bool = False, +) -> dict[str, int]: + """Import all messages from an mbox file. + + Returns stats: {processed, imported, skipped, errors}. + With dry_run=True counts messages without writing anything. + """ + stats: dict[str, int] = {"processed": 0, "imported": 0, "skipped": 0, "errors": 0} + + conn: Optional[asyncpg.Connection] = None + if not dry_run and dsn: + conn = await asyncpg.connect(dsn) + + try: + mbox = mailbox.mbox(str(mbox_path), create=False) + for msg in mbox: + stats["processed"] += 1 + if dry_run: + stats["imported"] += 1 + else: + outcome = await _import_one(msg, archive_root, conn) + if outcome == "imported": + stats["imported"] += 1 + elif outcome == "skipped": + stats["skipped"] += 1 + else: + stats["errors"] += 1 + if stats["processed"] % 500 == 0: + _log.info("progress", **stats) + finally: + if conn: + await conn.close() + + _log.info("import_complete", dry_run=dry_run, **stats) + return stats + + +def main() -> None: + parser = argparse.ArgumentParser( + description="Import Gmail Takeout .mbox into KB archive and envelope DB." + ) + parser.add_argument("--mbox", required=True, type=Path, metavar="PATH", + help="Path to Gmail Takeout .mbox file") + parser.add_argument("--archive", required=True, type=Path, metavar="DIR", + help="Archive root directory") + parser.add_argument("--dsn", metavar="DSN", + help="asyncpg DSN for kb-postgres (omit for archive-only mode)") + parser.add_argument("--dry-run", action="store_true", + help="Count messages only; no writes to archive or DB") + args = parser.parse_args() + + if not args.mbox.is_file(): + _log.error("mbox_not_found", path=str(args.mbox)) + sys.exit(1) + + stats = asyncio.run( + run_import(args.mbox, args.archive, args.dsn or None, args.dry_run) + ) + sys.exit(1 if stats["errors"] > 0 else 0) + + +if __name__ == "__main__": + main() diff --git a/jobs/gmail-bulk-import/tests/test_importer.py b/jobs/gmail-bulk-import/tests/test_importer.py new file mode 100644 index 0000000..e29d442 --- /dev/null +++ b/jobs/gmail-bulk-import/tests/test_importer.py @@ -0,0 +1,148 @@ +"""Unit tests for Gmail bulk importer — no DB, no external services required.""" +from __future__ import annotations + +import mailbox +from datetime import datetime, timezone +from pathlib import Path + +import pytest + +from gmail_bulk_import.importer import _message_id, _parse_date, run_import + + +def _make_msg(headers: dict) -> mailbox.mboxMessage: + msg = mailbox.mboxMessage() + for k, v in headers.items(): + msg[k] = v + return msg + + +def _make_mbox(path: Path, msgs: list[mailbox.mboxMessage]) -> None: + mbox = mailbox.mbox(str(path), create=True) + for msg in msgs: + mbox.add(msg) + mbox.flush() + mbox.close() + + +def _msg_with_id(mid: str = "") -> mailbox.mboxMessage: + return _make_msg({ + "From": "alice@example.com", + "Message-ID": mid, + "Date": "Tue, 10 Jun 2025 12:00:00 +0000", + "Subject": "Hello", + }) + + +class TestMessageId: + def test_uses_message_id_header(self): + msg = _make_msg({"Message-ID": ""}) + assert _message_id(msg) == "foo-123@example.com" + + def test_strips_angle_brackets(self): + msg = _make_msg({"Message-ID": " "}) + assert _message_id(msg) == "bar@baz.com" + + def test_fallback_to_sha256_when_missing(self): + msg = _make_msg({"From": "a@b.com", "Subject": "no id"}) + result = _message_id(msg) + assert result.startswith("sha256-") + assert len(result) == len("sha256-") + 32 + + def test_sha256_is_stable(self): + msg = _make_msg({"From": "a@b.com", "Subject": "stable"}) + assert _message_id(msg) == _message_id(msg) + + +class TestParseDate: + def test_parses_rfc2822(self): + msg = _make_msg({"Date": "Tue, 10 Jun 2025 12:00:00 +0000"}) + ts = _parse_date(msg) + assert ts == datetime(2025, 6, 10, 12, 0, 0, tzinfo=timezone.utc) + + def test_result_is_utc_aware(self): + msg = _make_msg({"Date": "Tue, 10 Jun 2025 12:00:00 +0000"}) + assert _parse_date(msg).tzinfo is not None + + def test_falls_back_to_epoch_on_bad_date(self): + msg = _make_msg({"Date": "not-a-date"}) + assert _parse_date(msg) == datetime(1970, 1, 1, tzinfo=timezone.utc) + + def test_falls_back_to_epoch_when_missing(self): + msg = _make_msg({"From": "a@b.com"}) + assert _parse_date(msg) == datetime(1970, 1, 1, tzinfo=timezone.utc) + + def test_converts_offset_to_utc(self): + msg = _make_msg({"Date": "Tue, 10 Jun 2025 14:00:00 +0200"}) + ts = _parse_date(msg) + assert ts == datetime(2025, 6, 10, 12, 0, 0, tzinfo=timezone.utc) + + +class TestRunImport: + async def test_dry_run_counts_without_writing(self, tmp_path): + mbox_path = tmp_path / "mail.mbox" + archive = tmp_path / "archive" + _make_mbox(mbox_path, [_msg_with_id(""), _msg_with_id("")]) + + stats = await run_import(mbox_path, archive, dry_run=True) + + assert stats["processed"] == 2 + assert stats["imported"] == 2 + assert not archive.exists() + + async def test_imports_eml_to_archive(self, tmp_path): + mbox_path = tmp_path / "mail.mbox" + archive = tmp_path / "archive" + _make_mbox(mbox_path, [_msg_with_id()]) + + stats = await run_import(mbox_path, archive) + + assert stats["imported"] == 1 + assert stats["skipped"] == 0 + assert stats["errors"] == 0 + eml_files = list(archive.rglob("*.eml")) + assert len(eml_files) == 1 + + async def test_eml_stored_under_gmail_source(self, tmp_path): + mbox_path = tmp_path / "mail.mbox" + archive = tmp_path / "archive" + _make_mbox(mbox_path, [_msg_with_id()]) + + await run_import(mbox_path, archive) + + eml_files = list(archive.rglob("*.eml")) + assert eml_files[0].parts[len(archive.parts)] == "gmail" + + async def test_idempotent_second_run_skips(self, tmp_path): + mbox_path = tmp_path / "mail.mbox" + archive = tmp_path / "archive" + _make_mbox(mbox_path, [_msg_with_id()]) + + await run_import(mbox_path, archive) + stats2 = await run_import(mbox_path, archive) + + assert stats2["imported"] == 0 + assert stats2["skipped"] == 1 + assert stats2["errors"] == 0 + + async def test_two_distinct_messages_both_imported(self, tmp_path): + mbox_path = tmp_path / "mail.mbox" + archive = tmp_path / "archive" + _make_mbox(mbox_path, [_msg_with_id(""), _msg_with_id("")]) + + stats = await run_import(mbox_path, archive) + + assert stats["imported"] == 2 + assert stats["errors"] == 0 + assert len(list(archive.rglob("*.eml"))) == 2 + + async def test_archive_only_mode_no_dsn(self, tmp_path): + """Without a DSN, job runs archive-only (no DB connection attempted).""" + mbox_path = tmp_path / "mail.mbox" + archive = tmp_path / "archive" + _make_mbox(mbox_path, [_msg_with_id()]) + + stats = await run_import(mbox_path, archive, dsn=None) + + assert stats["imported"] == 1 + assert list(archive.rglob("*.eml"))