Thicket/tests/test_pipeline_core.py

172 lines
6.8 KiB
Python

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