Thicket/thicket/pipeline_core.py

430 lines
18 KiB
Python

"""Pipeline orchestration — the shared engine behind GUI and CLI.
``IngestPipeline`` owns the per-file walk: extract, then a stage table
(vault note -> vector index -> knowledge graph). It is UI-agnostic:
progress flows out through ``PipelineCallbacks`` (plain callables),
which the QThread worker wires to Qt signals and the headless CLI
wires to prints. Heavy objects (embedding model, Qdrant client,
LightRAG) are constructed once inside ``run()`` — always on the
caller's thread, never the UI thread.
Design invariants:
* One stage table drives dispatch — adding a stage is one tuple,
never a new branch nest.
* Stages communicate through a per-file ``ctx`` dict (the vault
stage publishes the note's vault-relative path); no hidden
cross-stage state.
* A failing document is isolated: one ERROR row, queue continues.
* Stop is cooperative and checked between files only — a file is
either fully processed or not started.
Stage status vocabulary (the UI color-codes on these exact strings):
QUEUED, EXTRACT, VAULT, INDEX, GRAPH, DONE, SKIP, ERROR.
"""
from __future__ import annotations
import bz2
import hashlib
import json
import shutil
import sys
from collections.abc import Callable
from dataclasses import dataclass
from pathlib import Path
from .chunker import ContextualChunker
from .embedder import DEFAULT_EMBED_MODEL, EmbeddingEngine, model_dim
from .extractors import DocumentExtractor, scan_files
from .graph_store import create_graph
from .layout import default_archive, expand_path
from .minio_archive import MinioArchiver
from .vector_stores import TARGETS, VectorStoreError, create_store, doc_key
from .vault_writer import ObsidianVaultWriter
# Stage runner: (filepath, title, content, ctx) -> detail line.
StageRunner = Callable[[Path, str, str, dict], str]
VAULT_OUTPUT_DIRNAME = "Ingested_Brain"
# The notes-only destination key — valid wherever a target is accepted.
OBSIDIAN_TARGET = "obsidian"
def _bz2_compress(source: Path, level: int = 9) -> Path:
"""Streaming bzip2 of *source*; the plain original is replaced by
the .bz2 (compress-then-unlink, never unlink-first)."""
target = source.with_suffix(source.suffix + ".bz2")
with source.open("rb") as src, bz2.BZ2File(target, "wb",
compresslevel=level) as dst:
shutil.copyfileobj(src, dst, length=1 << 20)
source.unlink()
return target
@dataclass(slots=True)
class IngestConfig:
input_dir: Path
vault_dir: Path
qdrant_host: str = "localhost"
qdrant_port: int = 6333
target: str = "qdrant"
collection: str = "second_brain"
embed_model: str = DEFAULT_EMBED_MODEL
chunk_size: int = 400
overlap: int = 50
# Destination model: TARGET picks where the corpus lands — a
# vector store key from the registry, or "obsidian" for notes-only.
target: str = "qdrant" # kept in sync with qdrant_host below
use_vault: bool = True
use_qdrant: bool = True
use_lightrag: bool = False
use_minio: bool = False
use_fs_archive: bool = False
archive_dir: Path | None = None # default: corpus/archive or sibling
graph_engine: str = "lightrag"
notes_dir: str = VAULT_OUTPUT_DIRNAME
skip_unchanged: bool = False
max_mb: int = 0 # 0 = no size guard
def __post_init__(self) -> None:
"""Resolve ~ and $VAR in every user-supplied path — the GUI,
CLI, and core share one expansion point."""
self.input_dir = expand_path(self.input_dir)
self.vault_dir = expand_path(self.vault_dir)
if self.archive_dir is not None:
self.archive_dir = expand_path(self.archive_dir)
ollama_llm: str = "llama3"
ollama_embed: str = "nomic-embed-text"
@dataclass(slots=True)
class PipelineCallbacks:
log: object = lambda _msg: None # (msg)
file_status: object = lambda _p, _s: None # (path, status)
file_detail: object = lambda _p, _d: None # (path, detail)
progress: object = lambda _n, _c, _t: None # (name, current, total)
@classmethod
def quiet(cls) -> "PipelineCallbacks":
return cls()
class IngestPipeline:
"""One batch run over a directory of documents."""
def __init__(self, config: IngestConfig,
callbacks: PipelineCallbacks = PipelineCallbacks.quiet()):
self.config = config
self._manifest_path = config.vault_dir / ".thicket" / "manifest.json"
self.cb = callbacks
self._stop = False
self._engine: EmbeddingEngine | None = None
self._store = None
self._writer: ObsidianVaultWriter | None = None
self._graph = None
self._archiver = None
self._chunker = ContextualChunker(
chunk_size=config.chunk_size, overlap=config.overlap
)
# Stage table: (status, gate, runner). Dispatch is one pass
# over this tuple — stage order and gating live here only.
self._stages: tuple[tuple[str, bool, StageRunner], ...] = (
("MINIO", config.use_minio, self._stage_minio),
("VAULT", config.use_vault, self._stage_vault),
("INDEX", config.use_qdrant, self._stage_index),
("GRAPH", config.use_lightrag, self._stage_graph),
("ARCHIVE", config.use_fs_archive, self._stage_fs_archive),
)
def request_stop(self) -> None:
"""Cooperative stop: the queue exits after the current file."""
self._stop = True
# ── change manifest (powers skip-unchanged) ──
def _read_manifest(self) -> dict:
try:
return json.loads(self._manifest_path.read_text(encoding="utf-8"))
except (OSError, ValueError):
return {}
def _write_manifest(self, manifest: dict) -> None:
self._manifest_path.parent.mkdir(parents=True, exist_ok=True)
self._manifest_path.write_text(
json.dumps(manifest, indent=1), encoding="utf-8")
@staticmethod
def _file_digest(path: Path) -> str:
return hashlib.md5(path.read_bytes()).hexdigest()
# ── lazy stage construction (runs on the worker thread) ──
def _get_writer(self) -> ObsidianVaultWriter:
import importlib.util
if importlib.util.find_spec("slugify") is None:
raise VectorStoreError(
"python-slugify not installed — run: pip install 'thicket[ingest]'"
)
if self._writer is None:
self._writer = ObsidianVaultWriter(
self.config.vault_dir, output_dirname=self.config.notes_dir
)
return self._writer
def _get_engine(self) -> EmbeddingEngine:
if self._engine is None:
self.cb.log(
f"Loading local embedding model '{self.config.embed_model}' "
f"(first run downloads it)..."
)
self._engine = EmbeddingEngine(self.config.embed_model)
self._engine.load()
self.cb.log(f"Embedding model ready ({self._engine.dim} dimensions).")
return self._engine
def _get_store(self):
if self._store is None:
spec = TARGETS.get(self.config.target)
if spec is None:
known = ", ".join(sorted(TARGETS))
raise VectorStoreError(
f"unknown vector target '{self.config.target}' — known: {known}"
)
store = create_store(
self.config.target,
collection=self.config.collection,
dim=model_dim(self.config.embed_model),
host=self.config.qdrant_host,
port=self.config.qdrant_port,
data_dir=self._store_data_dir(),
log=self.cb.log,
)
store.set_embedder(self._get_engine())
store.ensure_collection()
self._store = store
return self._store
def _store_data_dir(self) -> Path:
"""Embedded/file targets keep data under the vault so the whole
knowledge tree stays one portable directory."""
return self.config.vault_dir / ".thicket" / self.config.target
def _get_graph(self):
if self._graph is None:
self.cb.log(
f"Initializing LightRAG graph at "
f"'{self.config.vault_dir / '.lightrag'}' "
f"(LLM: {self.config.ollama_llm}, embed: {self.config.ollama_embed})..."
)
self._graph = create_graph(
self.config.graph_engine,
working_dir=self.config.vault_dir,
llm_model=self.config.ollama_llm,
embed_model=self.config.ollama_embed,
log=self.cb.log,
)
return self._graph
def _get_archiver(self) -> MinioArchiver:
if self._archiver is None:
self._archiver = MinioArchiver(log=self.cb.log)
return self._archiver
# ── stage runners ──
def _stage_minio(self, filepath: Path, title: str, content: str,
ctx: dict) -> str:
uri = self._get_archiver().archive(filepath, self.config.input_dir)
ctx["source_uri"] = uri
self.cb.log(f" └─ Archived to MinIO: {uri}")
return f"s3: {uri}"
def _stage_fs_archive(self, filepath: Path, title: str, content: str,
ctx: dict) -> str:
"""Move the source out of the incoming tree into the archive
directory and bzip2 it — sources are never deleted."""
archive_dir = self.config.archive_dir or default_archive(
self.config.input_dir)
archive_dir.mkdir(parents=True, exist_ok=True)
target = archive_dir / filepath.name
if target.exists():
digest = doc_key(title, str(filepath))[:6]
target = archive_dir / f"{filepath.stem}-{digest}{filepath.suffix}"
shutil.move(str(filepath), target)
compressed = _bz2_compress(target)
self.cb.log(f" └─ Archived + bz2: {compressed.relative_to(compressed.parents[1])}")
return f"bz2: {compressed.name}"
def _stage_vault(self, filepath: Path, title: str, content: str,
ctx: dict) -> str:
note = self._get_writer().write(
title, content, filepath, source_uri=ctx.get("source_uri"))
ctx["rel_path"] = str(note.relative_to(self.config.vault_dir))
self.cb.log(f" └─ Vault note written: [[{note.stem}]]")
return f"note: {note.name}"
def _stage_index(self, filepath: Path, title: str, content: str,
ctx: dict) -> str:
rel_path = ctx.get("rel_path") or self._projected_rel_path(title)
chunks = self._chunker.chunk(
title, content,
source_path=str(filepath.relative_to(self.config.input_dir)))
count = self._get_store().replace_document(title, rel_path, chunks)
self.cb.log(
f" └─ Vector store: indexed {count} chunk(s) into "
f"'{self.config.collection}'."
)
return f"{count} chunks indexed"
def _stage_graph(self, filepath: Path, title: str, content: str,
ctx: dict) -> str:
self._get_graph().ingest_document(title, content)
self.cb.log(" └─ LightRAG graph updated.")
return "graph updated"
def _projected_rel_path(self, title: str) -> str:
"""Canonical note path for the vector payload when the vault
stage is disabled — identical location to the writer's plain
slug (collision suffixes only exist once notes are written)."""
from slugify import slugify
slug = slugify(title) or "untitled"
return f"{self.config.notes_dir}/{slug}.md"
# ── main loop ──
def run(self, files: list[Path] | None = None) -> tuple[int, int]:
"""Process every supported document. Returns (ok, fail)."""
files = scan_files(self.config.input_dir) if files is None else files
self.cb.log(f"Found {len(files)} eligible document(s) in "
f"'{self.config.input_dir}'")
for f in files:
self.cb.file_status(str(f), "QUEUED")
outcomes = [self._process_one(f) for f in self._itinerary(files)]
ok = sum(outcomes)
if self._graph is not None:
try:
self._graph.finalize() # batch engines build once here
except Exception as e: # noqa: BLE001 — report, keep counts
self.cb.log(f" └─ ERROR finalizing graph: {e}")
return ok, len(outcomes) - ok
def _itinerary(self, files: list[Path]):
"""Yield files in order, honouring the stop flag between files
and emitting progress as we go."""
total = len(files)
for current, filepath in enumerate(files, start=1):
if self._stop:
self.cb.log("STOP: exiting queue before next file.")
return
self.cb.progress(filepath.name, current, total)
yield filepath
def _process_one(self, filepath: Path) -> bool:
"""One document through the gates + extract + the stage table."""
try:
skip_reason = self._gate(filepath)
if skip_reason is not None:
self.cb.file_status(str(filepath), "SKIP")
self.cb.file_detail(str(filepath), skip_reason)
self.cb.log(f" └─ Skipping {filepath.name}: {skip_reason}")
return True # deliberate non-processing — not a failure
self.cb.file_status(str(filepath), "EXTRACT")
self.cb.file_detail(str(filepath), "extracting text")
title, content = DocumentExtractor.extract(filepath)
if not content.strip():
self.cb.log(f" └─ Skipping empty file: {filepath.name}")
self.cb.file_status(str(filepath), "SKIP")
self.cb.file_detail(str(filepath), "no extractable text")
return True # nothing to ingest — not a failure
ctx: dict = {}
for status, enabled, runner in self._stages:
if not enabled:
continue
self.cb.file_status(str(filepath), status)
self.cb.file_detail(str(filepath), runner(filepath, title, content, ctx))
self.cb.file_status(str(filepath), "DONE")
if self.config.skip_unchanged:
manifest = self._read_manifest()
manifest[str(filepath)] = self._file_digest(filepath)
self._write_manifest(manifest)
return True
except Exception as e: # noqa: BLE001 — per-file isolation is the contract
self._fail(filepath, e)
return False
def _gate(self, filepath: Path) -> str | None:
"""Pre-extraction gates: oversized files and unchanged re-runs.
Returns the skip reason, or None to proceed."""
if self.config.max_mb:
size_mb = filepath.stat().st_size / 1_048_576
if size_mb > self.config.max_mb:
return f"{size_mb:.1f} MB exceeds {self.config.max_mb} MB limit"
if self.config.skip_unchanged:
manifest = self._read_manifest()
if manifest.get(str(filepath)) == self._file_digest(filepath):
return "unchanged since last ingest"
return None
def _fail(self, filepath: Path, error: Exception) -> None:
msg = f" └─ ERROR processing {filepath.name}: {error}"
print(msg, file=sys.stderr)
self.cb.log(msg)
self.cb.file_status(str(filepath), "ERROR")
self.cb.file_detail(str(filepath), str(error)[:120])
# ──────────────────────────────────────────────────────────────────
# HEADLESS ENTRY (no Qt)
# ──────────────────────────────────────────────────────────────────
def headless_ingest(input_dir: Path, vault_dir: Path,
qdrant_host: str = "localhost", qdrant_port: int = 6333,
target: str = "qdrant", collection: str = "second_brain",
embed_model: str = DEFAULT_EMBED_MODEL,
chunk_size: int = 400, overlap: int = 50,
use_vault: bool = True, use_qdrant: bool = True,
use_lightrag: bool = False, use_minio: bool = False,
graph_engine: str = "lightrag",
notes_dir: str = VAULT_OUTPUT_DIRNAME,
skip_unchanged: bool = False,
use_fs_archive: bool = False,
archive_dir: Path | None = None, max_mb: int = 0,
ollama_llm: str = "llama3",
ollama_embed: str = "nomic-embed-text") -> tuple[int, int]:
"""CLI batch ingest — the same core path the GUI worker drives."""
config = IngestConfig(
input_dir=input_dir, vault_dir=vault_dir,
qdrant_host=qdrant_host, qdrant_port=qdrant_port,
target=target, collection=collection, embed_model=embed_model,
chunk_size=chunk_size, overlap=overlap,
use_vault=use_vault, use_qdrant=use_qdrant, use_lightrag=use_lightrag,
use_minio=use_minio, graph_engine=graph_engine,
notes_dir=notes_dir, skip_unchanged=skip_unchanged,
use_fs_archive=use_fs_archive, archive_dir=archive_dir,
max_mb=max_mb,
ollama_llm=ollama_llm, ollama_embed=ollama_embed,
)
callbacks = PipelineCallbacks(
log=lambda msg: print(f"> {msg}"),
file_status=lambda _p, _s: None,
file_detail=lambda _p, d: print(f" {d}"),
progress=lambda n, c, t: print(f"\n[{c}/{t}] {n}"),
)
pipeline = IngestPipeline(config, callbacks)
try:
return pipeline.run()
except KeyboardInterrupt:
# POSIX: exit code 130 = 128 + SIGINT; report what completed.
print("\nInterrupted — counts reflect files finished before SIGINT.")
raise