430 lines
18 KiB
Python
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
|