172 lines
6.8 KiB
Python
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()
|