"""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