Cognitive-rag / 01_data_ingestion / tests / test_integration.py
test_integration.py
Raw
"""Integration / smoke tests that run the full ingestion pipeline on small test files."""

from __future__ import annotations

import json
import sys
from pathlib import Path

import pytest

sys.path.insert(0, str(Path(__file__).resolve().parent.parent))

from ingestion.engine import IngestionConfig, IngestionEngine
from ingestion.events import CallbackEmitter
from tests.conftest import read_jsonl


# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------

def _make_config(
    input_dir: Path,
    output_dir: Path,
    workers: int = 2,
) -> IngestionConfig:
    return IngestionConfig(
        input_dir=input_dir,
        output_dir=output_dir,
        workers=workers,
    )


# ---------------------------------------------------------------------------
# End-to-end: text + CSV
# ---------------------------------------------------------------------------

class TestTextAndCSVIngestion:
    def test_ingests_txt_and_csv(self, tmp_path: Path):
        input_dir = tmp_path / "input"
        input_dir.mkdir()
        output_dir = tmp_path / "output"

        (input_dir / "readme.txt").write_text(
            "This is a test document with some content.\n"
            "It has multiple lines of text."
        )
        (input_dir / "data.csv").write_text(
            "name,age,city\n"
            "Alice,30,NYC\n"
            "Bob,25,LA\n"
            "Charlie,35,Chicago\n"
        )

        config = _make_config(input_dir, output_dir)
        stats = IngestionEngine(config).run()

        assert stats["counts"]["files_succeeded"] == 2
        assert stats["counts"]["files_failed"] == 0
        assert stats["counts"]["records_emitted"] >= 4  # 1 txt + 3 csv rows

        records = read_jsonl(output_dir / "cleaned_documents.jsonl")
        assert len(records) >= 4

        # Check record schema
        for rec in records:
            assert "doc_id" in rec
            assert "content" in rec
            assert "source" in rec
            assert "record_id" in rec["source"]
            assert rec["content"].strip() != ""

    def test_csv_preserves_column_names(self, tmp_path: Path):
        input_dir = tmp_path / "input"
        input_dir.mkdir()
        output_dir = tmp_path / "output"

        (input_dir / "test.csv").write_text("product,price\nWidget,9.99\nGadget,19.99\n")

        config = _make_config(input_dir, output_dir)
        IngestionEngine(config).run()

        records = read_jsonl(output_dir / "cleaned_documents.jsonl")
        assert len(records) == 2
        assert "product" in records[0]["content"]
        assert "Widget" in records[0]["content"]


# ---------------------------------------------------------------------------
# Nested directories
# ---------------------------------------------------------------------------

class TestNestedDirectories:
    def test_recursive_ingestion(self, tmp_path: Path):
        input_dir = tmp_path / "input"
        input_dir.mkdir()
        output_dir = tmp_path / "output"

        # Create nested structure
        (input_dir / "top.txt").write_text("Top level content")
        sub1 = input_dir / "level1"
        sub1.mkdir()
        (sub1 / "mid.txt").write_text("Mid level content")
        sub2 = sub1 / "level2"
        sub2.mkdir()
        (sub2 / "deep.txt").write_text("Deep level content")

        config = _make_config(input_dir, output_dir)
        stats = IngestionEngine(config).run()

        assert stats["counts"]["files_succeeded"] == 3
        records = read_jsonl(output_dir / "cleaned_documents.jsonl")
        contents = [r["content"] for r in records]
        assert any("Top level" in c for c in contents)
        assert any("Mid level" in c for c in contents)
        assert any("Deep level" in c for c in contents)

    def test_nested_paths_in_source(self, tmp_path: Path):
        input_dir = tmp_path / "input"
        input_dir.mkdir()
        output_dir = tmp_path / "output"

        sub = input_dir / "docs" / "legal"
        sub.mkdir(parents=True)
        (sub / "policy.txt").write_text("Policy content")

        config = _make_config(input_dir, output_dir)
        IngestionEngine(config).run()

        records = read_jsonl(output_dir / "cleaned_documents.jsonl")
        assert len(records) == 1
        assert "docs/legal/policy.txt" in records[0]["source"]["path"]


# ---------------------------------------------------------------------------
# Unsupported file reporting
# ---------------------------------------------------------------------------

class TestUnsupportedFiles:
    def test_skipped_files_reported(self, tmp_path: Path):
        input_dir = tmp_path / "input"
        input_dir.mkdir()
        output_dir = tmp_path / "output"

        (input_dir / "good.txt").write_text("Valid content")
        (input_dir / "program.exe").write_bytes(b"\x00\x01\x02\x03")
        (input_dir / "library.dll").write_bytes(b"\x4d\x5a")
        (input_dir / "data.bin").write_bytes(b"\xff\xfe")

        received = []
        emitter = CallbackEmitter(received.append)
        config = _make_config(input_dir, output_dir)
        stats = IngestionEngine(config, emitter=emitter).run()

        assert stats["counts"]["files_skipped"] == 3
        assert stats["counts"]["files_succeeded"] == 1

        skip_events = [e for e in received if e.event_type == "file_skipped"]
        assert len(skip_events) == 3

        skipped_paths = {e.file_path for e in skip_events}
        assert any("exe" in p for p in skipped_paths)
        assert any("dll" in p for p in skipped_paths)
        assert any("bin" in p for p in skipped_paths)

    def test_skipped_files_in_stats_json(self, tmp_path: Path):
        input_dir = tmp_path / "input"
        input_dir.mkdir()
        output_dir = tmp_path / "output"

        (input_dir / "doc.txt").write_text("Content")
        (input_dir / "app.exe").write_bytes(b"\x00")

        config = _make_config(input_dir, output_dir)
        IngestionEngine(config).run()

        with open(output_dir / "ingestion_stats.json") as f:
            stats = json.load(f)
        assert len(stats["skipped_files"]) == 1
        assert "exe" in stats["skipped_files"][0]["reason"].lower() or "exe" in stats["skipped_files"][0]["path"]


# ---------------------------------------------------------------------------
# Record relations (prev/next)
# ---------------------------------------------------------------------------

class TestRecordRelations:
    def test_csv_records_have_relations(self, tmp_path: Path):
        input_dir = tmp_path / "input"
        input_dir.mkdir()
        output_dir = tmp_path / "output"

        (input_dir / "data.csv").write_text("a,b\n1,2\n3,4\n5,6\n")

        config = _make_config(input_dir, output_dir)
        IngestionEngine(config).run()

        records = read_jsonl(output_dir / "cleaned_documents.jsonl")
        assert len(records) == 3

        # First record: no prev, has next
        assert records[0]["source"]["prev_record_id"] is None
        assert records[0]["source"]["next_record_id"] is not None

        # Middle record: has both
        assert records[1]["source"]["prev_record_id"] is not None
        assert records[1]["source"]["next_record_id"] is not None

        # Last record: has prev, no next
        assert records[2]["source"]["prev_record_id"] is not None
        assert records[2]["source"]["next_record_id"] is None


# ---------------------------------------------------------------------------
# Deterministic output
# ---------------------------------------------------------------------------

class TestDeterministicOutput:
    def test_same_input_same_output(self, tmp_path: Path):
        input_dir = tmp_path / "input"
        input_dir.mkdir()
        (input_dir / "a.txt").write_text("Alpha")
        (input_dir / "b.txt").write_text("Beta")
        (input_dir / "c.csv").write_text("x\n1\n2\n")

        out1 = tmp_path / "out1"
        out2 = tmp_path / "out2"

        cfg1 = _make_config(input_dir, out1, workers=4)
        IngestionEngine(cfg1).run()
        cfg2 = _make_config(input_dir, out2, workers=4)
        IngestionEngine(cfg2).run()

        recs1 = read_jsonl(out1 / "cleaned_documents.jsonl")
        recs2 = read_jsonl(out2 / "cleaned_documents.jsonl")

        assert len(recs1) == len(recs2)
        for r1, r2 in zip(recs1, recs2):
            assert r1["doc_id"] == r2["doc_id"]
            assert r1["content"] == r2["content"]


# ---------------------------------------------------------------------------
# Empty input
# ---------------------------------------------------------------------------

class TestEmptyInput:
    def test_empty_directory(self, tmp_path: Path):
        input_dir = tmp_path / "empty_input"
        input_dir.mkdir()
        output_dir = tmp_path / "output"

        config = _make_config(input_dir, output_dir)
        stats = IngestionEngine(config).run()

        assert stats["counts"]["files_processed"] == 0
        assert stats["counts"]["records_emitted"] == 0

        output_path = output_dir / "cleaned_documents.jsonl"
        assert output_path.exists()
        assert output_path.read_text().strip() == ""

    def test_only_unsupported_files(self, tmp_path: Path):
        input_dir = tmp_path / "input"
        input_dir.mkdir()
        output_dir = tmp_path / "output"

        (input_dir / "a.exe").write_bytes(b"\x00")
        (input_dir / "b.dll").write_bytes(b"\x00")

        config = _make_config(input_dir, output_dir)
        stats = IngestionEngine(config).run()

        assert stats["counts"]["files_processed"] == 0
        assert stats["counts"]["files_skipped"] == 2