"""IngestPipeline — vault-only end-to-end run (no services needed).""" from __future__ import annotations from thicket.pipeline_core import ( IngestConfig, IngestPipeline, PipelineCallbacks, ) class RecordingCallbacks: def __init__(self): self.logs: list[str] = [] self.statuses: list[tuple[str, str]] = [] self.details: list[tuple[str, str]] = [] self.progress: list[tuple[str, int, int]] = [] def bind(self) -> PipelineCallbacks: return PipelineCallbacks( log=self.logs.append, file_status=lambda p, s: self.statuses.append((p, s)), file_detail=lambda p, d: self.details.append((p, d)), progress=lambda n, c, t: self.progress.append((n, c, t)), ) def test_vault_only_pipeline_processes_every_document(sample_docs, vault): config = IngestConfig( input_dir=sample_docs, vault_dir=vault, use_qdrant=False, use_lightrag=False, ) cb = RecordingCallbacks() ok, fail = IngestPipeline(config, cb.bind()).run() assert (ok, fail) == (2, 0) assert (vault / "Ingested_Brain" / "deep-learning-notes.md").exists() assert (vault / "Ingested_Brain" / "reading-list.md").exists() # every file walked the full status lifecycle for path in {p for p, _ in cb.statuses}: stages = [s for p, s in cb.statuses if p == path] assert stages[0] == "QUEUED" assert stages[-1] == "DONE" assert "EXTRACT" in stages and "VAULT" in stages # progress counted 1..2 assert cb.progress[-1][1:] == (2, 2) def test_empty_file_is_skip_not_failure(sample_docs, vault): (sample_docs / "blank.md").write_text("", encoding="utf-8") config = IngestConfig(input_dir=sample_docs, vault_dir=vault, use_qdrant=False, use_lightrag=False) cb = RecordingCallbacks() ok, fail = IngestPipeline(config, cb.bind()).run() assert (ok, fail) == (3, 0) assert ("SKIP" in [s for _, s in cb.statuses]) def test_cooperative_stop_between_files(sample_docs, vault): config = IngestConfig(input_dir=sample_docs, vault_dir=vault, use_qdrant=False, use_lightrag=False) pipeline = IngestPipeline(config, PipelineCallbacks.quiet()) original = pipeline._process_one calls = {"n": 0} def stop_after_first(filepath): calls["n"] += 1 result = original(filepath) pipeline.request_stop() return result pipeline._process_one = stop_after_first ok, fail = pipeline.run() assert calls["n"] == 1 # stopped before the second file assert (ok, fail) == (1, 0) def test_broken_document_fails_alone_and_queue_continues(sample_docs, vault): (sample_docs / "broken.epub").write_bytes(b"not really an epub") config = IngestConfig(input_dir=sample_docs, vault_dir=vault, use_qdrant=False, use_lightrag=False) cb = RecordingCallbacks() ok, fail = IngestPipeline(config, cb.bind()).run() # the two good files still succeeded; the epub failed in isolation assert ok == 2 and fail == 1 assert any("ERROR processing broken.epub" in line for line in cb.logs) def test_vault_stage_disabled_writes_no_notes(sample_docs, vault): config = IngestConfig(input_dir=sample_docs, vault_dir=vault, use_vault=False, use_qdrant=False, use_lightrag=False) cb = RecordingCallbacks() ok, fail = IngestPipeline(config, cb.bind()).run() assert (ok, fail) == (2, 0) assert not (vault / "Ingested_Brain").exists() or \ not any((vault / "Ingested_Brain").iterdir()) assert not any("Vault note written" in line for line in cb.logs) # VAULT never appears in the stage lifecycle assert "VAULT" not in [s for _, s in cb.statuses] def test_skip_unchanged_skips_second_run(sample_docs, vault): config = IngestConfig(input_dir=sample_docs, vault_dir=vault, use_qdrant=False, use_lightrag=False, skip_unchanged=True) cb = RecordingCallbacks() first = IngestPipeline(config, cb.bind()).run() second = IngestPipeline(config, cb.bind()).run() assert first == (2, 0) assert second == (2, 0) assert [s for _, s in cb.statuses].count("SKIP") == 2 assert any("unchanged" in d for _, d in cb.details) # Editing a file brings it back into the queue. (sample_docs / "reading-list.txt").write_text("new content\n", encoding="utf-8") IngestPipeline(config, cb.bind()).run() assert [s for _, s in cb.statuses].count("DONE") == 3 def test_fs_archive_moves_and_compresses(sample_docs, vault, tmp_path): archive = tmp_path / "ingested-archive" config = IngestConfig(input_dir=sample_docs, vault_dir=vault, use_qdrant=False, use_lightrag=False, use_fs_archive=True, archive_dir=archive) ok, fail = IngestPipeline(config, PipelineCallbacks.quiet()).run() assert (ok, fail) == (2, 0) # Incoming tree is empty; the archive holds bz2 payloads only. assert not any(p.exists() for p in sample_docs.iterdir() if p.is_file()) packed = sorted(archive.glob("*.bz2")) assert len(packed) == 2 and not any(p.suffix != ".bz2" for p in archive.iterdir()) # Content round-trips through bunzip2. import bz2 text = bz2.decompress(packed[0].read_bytes()).decode("utf-8") assert text.strip() # real content survived def test_obsidian_target_is_notes_only(sample_docs, vault): config = IngestConfig(input_dir=sample_docs, vault_dir=vault, target="obsidian", use_qdrant=False, use_lightrag=False) cb = RecordingCallbacks() ok, fail = IngestPipeline(config, cb.bind()).run() assert (ok, fail) == (2, 0) assert (vault / "Ingested_Brain" / "reading-list.md").exists() assert "INDEX" not in [s_ for _, s_ in cb.statuses] # no vector stage def test_max_mb_guards_oversized_files(sample_docs, vault): (sample_docs / "huge.txt").write_text("x" * 1_200_000, encoding="utf-8") config = IngestConfig(input_dir=sample_docs, vault_dir=vault, use_qdrant=False, use_lightrag=False, max_mb=1) cb = RecordingCallbacks() ok, fail = IngestPipeline(config, cb.bind()).run() assert (ok, fail) == (3, 0) details = [d for _, d in cb.details] assert any("exceeds" in d for d in details) assert (sample_docs / "huge.txt").exists() # guards skip, never delete def test_custom_notes_dir(sample_docs, vault): config = IngestConfig(input_dir=sample_docs, vault_dir=vault, use_qdrant=False, use_lightrag=False, notes_dir="Technical_Shots") IngestPipeline(config, PipelineCallbacks.quiet()).run() assert (vault / "Technical_Shots" / "reading-list.md").exists()