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