import asyncio import json import logging import os import re import threading from collections.abc import Callable from datetime import datetime, timezone from pathlib import Path from typing import Any import frontmatter from backend.media_types import AUDIO_EXTENSIONS, IMAGE_EXTENSIONS, VIDEO_EXTENSIONS, is_media logger = logging.getLogger("obsigate.indexer") # Global in-memory index index: dict[str, dict[str, Any]] = {} # Vault config: {name: {path, attachmentsPath, scanAttachmentsOnStartup}} vault_config: dict[str, dict[str, Any]] = {} # Thread-safe lock for index updates _index_lock = threading.Lock() # Async lock for partial index updates (coexists with threading lock) _async_index_lock: asyncio.Lock | None = None # initialized lazily # Generation counter — incremented on each index rebuild so consumers # (e.g. the inverted index in search.py) can detect staleness. _index_generation: int = 0 # Timestamp of last full index rebuild (ISO format, empty if never built) _last_full_index_ts: str = "" # Hook for incremental inverted index updates: called as (action, vault, path, file_info) _on_index_change: Callable[..., None] | None = None def set_index_change_hook(hook): """Register a callback for incremental inverted index updates. The hook is called as ``hook(action, vault_name, path, file_info)`` where ``action`` is ``'add'`` or ``'remove'``. """ global _on_index_change _on_index_change = hook # O(1) lookup table for wikilink resolution: {filename_lower: [{vault, path}, ...]} _file_lookup: dict[str, list[dict[str, str]]] = {} # Backlink index: {vault_name: {relative_path: [{vault, path, title}, ...]}} _backlink_index: dict[str, dict[str, list[dict[str, str]]]] = {} # O(1) path index for tree filtering: {vault_name: [{path, name, type}, ...]} path_index: dict[str, list[dict[str, str]]] = {} # Maximum content size stored per file for in-memory search (bytes) SEARCH_CONTENT_LIMIT = 100_000 # Supported text-based file extensions SUPPORTED_EXTENSIONS = { ".md", ".txt", ".log", ".py", ".js", ".ts", ".jsx", ".tsx", ".sh", ".bash", ".zsh", ".fish", ".bat", ".cmd", ".ps1", ".json", ".yaml", ".yml", ".toml", ".xml", ".csv", ".cfg", ".ini", ".conf", ".env", ".pdf", ".xlsx", ".html", ".css", ".scss", ".less", ".java", ".c", ".cpp", ".h", ".hpp", ".cs", ".go", ".rs", ".rb", ".php", ".sql", ".r", ".m", ".swift", ".kt", ".dockerfile", ".makefile", ".cmake", ".excalidraw", ".excalidraw.md", } | set(IMAGE_EXTENSIONS) | set(AUDIO_EXTENSIONS) | set(VIDEO_EXTENSIONS) # Ignored directories (configurable via OBSIGATE_IGNORED_DIRS env var) _DEFAULT_IGNORED = {'.obsidian', '.trash', '.git', '__pycache__', 'node_modules', '.obsigate-backup'} _env_ignored = os.environ.get("OBSIGATE_IGNORED_DIRS", "") IGNORED_DIRS = {d.strip() for d in _env_ignored.split(",") if d.strip()} if _env_ignored else _DEFAULT_IGNORED.copy() def load_vault_config() -> dict[str, dict[str, Any]]: """Read VAULT_N_* and DIR_N_* env vars and return vault configuration. Scans environment variables ``VAULT_1_NAME``/``VAULT_1_PATH``, ``VAULT_2_NAME``/``VAULT_2_PATH``, etc. in sequential order. Stops at the first missing pair. Also reads optional configuration: - VAULT_N_ATTACHMENTS_PATH: relative path to attachments folder - VAULT_N_SCAN_ATTACHMENTS: "true"/"false" to enable/disable scanning Returns: Dict mapping vault names to configuration dicts with keys: - path: filesystem path (required) - attachmentsPath: relative attachments folder (optional) - scanAttachmentsOnStartup: boolean (default True) - type: "VAULT" or "DIR" """ vaults: dict[str, dict[str, Any]] = {} n = 1 while True: name = os.environ.get(f"VAULT_{n}_NAME") path = os.environ.get(f"VAULT_{n}_PATH") if not name or not path: break # Optional configuration attachments_path = os.environ.get(f"VAULT_{n}_ATTACHMENTS_PATH") scan_attachments = os.environ.get(f"VAULT_{n}_SCAN_ATTACHMENTS", "true").lower() == "true" vaults[name] = { "path": path, "attachmentsPath": attachments_path, "scanAttachmentsOnStartup": scan_attachments, "type": "VAULT" } n += 1 n = 1 while True: name = os.environ.get(f"DIR_{n}_NAME") path = os.environ.get(f"DIR_{n}_PATH") if not name or not path: break vaults[name] = { "path": path, "attachmentsPath": None, "scanAttachmentsOnStartup": False, "type": "DIR" } n += 1 return vaults # Regex for extracting inline #tags from markdown body (excludes code blocks) _INLINE_TAG_RE = re.compile(r'(?:^|\s)#([a-zA-Z][a-zA-Z0-9_/-]{1,50})', re.MULTILINE) # Regex patterns for stripping code blocks before inline tag extraction _CODE_BLOCK_RE = re.compile(r'```[\s\S]*?```', re.MULTILINE) _INLINE_CODE_RE = re.compile(r'`[^`]+`') def _extract_tags(post: frontmatter.Post) -> list[str]: """Extract tags from frontmatter metadata. Handles tags as comma-separated string, list, or other types. Strips leading ``#`` from each tag. Args: post: Parsed frontmatter Post object. Returns: List of cleaned tag strings. """ tags = post.metadata.get("tags", []) if isinstance(tags, str): tags = [t.strip().lstrip("#") for t in tags.split(",") if t.strip()] elif isinstance(tags, list): tags = [str(t).strip().lstrip("#") for t in tags] else: tags = [] return tags def _extract_inline_tags(content: str) -> list[str]: """Extract inline #tag patterns from markdown content. Strips fenced and inline code blocks before scanning to avoid false positives from code comments or shell commands. Args: content: Raw markdown content (without frontmatter). Returns: Deduplicated list of inline tag strings. """ stripped = _CODE_BLOCK_RE.sub('', content) stripped = _INLINE_CODE_RE.sub('', stripped) return list(set(_INLINE_TAG_RE.findall(stripped))) def _extract_title(post: frontmatter.Post, filepath: Path) -> str: """Extract title from frontmatter or derive from filename. Falls back to the file stem with hyphens/underscores replaced by spaces when no ``title`` key is present in frontmatter. Args: post: Parsed frontmatter Post object. filepath: Path to the source file. Returns: Human-readable title string. """ title = post.metadata.get("title", "") if not title: title = filepath.stem.replace("-", " ").replace("_", " ") return str(title) _EXCALIDRAW_TEXT_FIELDS = ("text", "originalText", "label", "title") def extract_excalidraw_text_from_elements(elements: list[dict[str, Any]]) -> str: """Concatenate all user-visible text from an Excalidraw elements list. Iterates diagram elements and collects the text-bearing fields (``text`` for text elements, ``label``/``title`` for bound shapes, ``originalText`` as the stable source of a text element). Non-text elements and geometry-only shapes contribute nothing. This gives the TF-IDF search engine human-readable content instead of raw JSON. """ chunks: list[str] = [] for el in elements: if not isinstance(el, dict): continue collected = set() for field in _EXCALIDRAW_TEXT_FIELDS: val = el.get(field) if isinstance(val, str) and val.strip(): collected.add(val.strip()) if collected: chunks.append(" ".join(sorted(collected))) return "\n".join(chunks) _EXCALIDRAW_B64_ALPHABET = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/=" def _decompress_excalidraw(compressed: str) -> dict[str, Any] | None: """Decompress the Obsidian ``compressed-json`` block of a .excalidraw.md file. The Obsidian Excalidraw plugin stores scene data as an ``lz-string`` ``compressToBase64`` payload (LZ-based compression, alphabet = standard base64). We port the reference ``lz-string`` ``_decompress`` algorithm (bitsPerChar=6/base64 key, resetValue=32) in pure Python so indexation can recover the diagram text without the JS client. The port is validated against multiple JS-truth fixtures in ``test_excalidraw.py`` including accented text and edge cases. Returns the parsed scene dict, or None if the payload is not valid base64-lz or does not parse as JSON. Only ``decompress`` is needed for indexation (read side); compression stays on the JS client. """ if not compressed: return None length = len(compressed) def _get_char_value(index: int) -> int: # Mirror JS: input.charAt(index), 0-indexed. ch = compressed[index] if 0 <= index < length else "=" pos = _EXCALIDRAW_B64_ALPHABET.find(ch) return 0 if pos == -1 else pos # ── Build a bit reader over the base64 characters ──────────────────── # JS decompressFromBase64 calls _decompress(length, 32, getNextValue). # state = current 6-bit value + position mask + next char index. value = _get_char_value(0) position = 32 # resetValue index = 1 def _read_bit() -> int: nonlocal value, position, index resb = value & position position >>= 1 if position == 0: position = 32 value = _get_char_value(index) index += 1 return 1 if resb > 0 else 0 def _read_bits(n: int) -> int: out = 0 for i in range(n): out |= _read_bit() << i return out # ── Decompressor state ──────────────────────────────────────────────── dictionary: list[Any] = [0, 1, 2] # values 0,1,2 as in JS enlarge_in = 4 dict_size = 4 num_bits = 3 first = _read_bits(2) if first == 0: c = chr(_read_bits(8)) elif first == 1: c = chr(_read_bits(16)) elif first == 2: return None # empty stream else: return None dictionary.append(c) w: str = c result = [c] while True: if index > length: return None code = _read_bits(num_bits) if code == 0: dictionary.append(chr(_read_bits(8))) dict_size += 1 code = dict_size - 1 enlarge_in -= 1 elif code == 1: dictionary.append(chr(_read_bits(16))) dict_size += 1 code = dict_size - 1 enlarge_in -= 1 elif code == 2: break # end of stream if enlarge_in == 0: enlarge_in = 1 << num_bits num_bits += 1 if code < len(dictionary) and dictionary[code]: entry = dictionary[code] elif code == dict_size: entry = w + w[0] else: return None result.append(entry) dictionary.append(w + entry[0]) dict_size += 1 enlarge_in -= 1 w = entry if enlarge_in == 0: enlarge_in = 1 << num_bits num_bits += 1 text = "".join(result) try: data = json.loads(text) except Exception: return None if not isinstance(data, dict): return None return data def extract_xlsx_indexable(file_path: Path) -> str: """Return searchable text for a workbook (#153 A5). Lazy wrapper: ``openpyxl`` is only imported when a spreadsheet is actually indexed, so a vault without workbooks never pays the import. Errors are swallowed — a corrupt or encrypted file still gets indexed by name. """ try: from backend.xlsx_reader import extract_indexable_text except Exception: # pragma: no cover - openpyxl missing return "" try: return extract_indexable_text(file_path) except Exception: # pragma: no cover - defensive return "" def extract_excalidraw_indexable(raw: str) -> str: """Return indexable text content for a raw .excalidraw / .excalidraw.md file. Supports both the pure JSON format (``type: "excalidraw"``) and the Obsidian ``compressed-json`` block. Falls back to ``""`` when neither format can be parsed so the file is still indexed (title only). """ # .excalidraw.md — Obsidian plugin embeds a compressed-json block if "excalidraw-plugin:" in raw: m = re.search(r"```compressed-json\n(.*?)\n```", raw, re.DOTALL) if m: parsed = _decompress_excalidraw(m.group(1).strip()) if parsed: return extract_excalidraw_text_from_elements(parsed.get("elements") or []) return "" # Pure .excalidraw JSON try: data = json.loads(raw) except Exception: return "" if not isinstance(data, dict) or data.get("type") != "excalidraw": return "" return extract_excalidraw_text_from_elements(data.get("elements") or []) def parse_markdown_file(raw: str) -> frontmatter.Post: """Parse markdown frontmatter, falling back to plain content if YAML is invalid. When the YAML block is malformed, strips it and returns a Post with empty metadata so that rendering can still proceed. Args: raw: Full raw markdown string including optional frontmatter. Returns: ``frontmatter.Post`` with ``.content`` and ``.metadata`` attributes. """ try: return frontmatter.loads(raw) except Exception as exc: logger.debug(f"Invalid frontmatter detected, falling back to plain markdown parsing: {exc}") content = raw if raw.startswith("---"): match = re.match(r"^---\s*\r?\n.*?\r?\n---\s*\r?\n?", raw, flags=re.DOTALL) if match: content = raw[match.end():] return frontmatter.Post(content) def _scan_vault( vault_name: str, vault_path: str, vault_cfg: dict[str, Any] | None = None, previous_files: dict[str, dict[str, Any]] | None = None, ) -> dict[str, Any]: """Synchronously scan a single vault directory and build file index. Walks the vault tree, reads supported files, extracts metadata (tags, title, content preview) and stores a capped content snapshot for in-memory full-text search. All files and directories are indexed, including hidden files (starting with '.'). Differential scan (#86): when ``previous_files`` maps a relative path to its previous ``file_info`` dict, entries whose ``size`` and ``modified`` timestamp are unchanged are reused verbatim (no disk read, no re-parse). Only the cheap ``os.walk`` + ``stat`` runs on every pass; heavy content extraction (PDF metadata excepted — always cheap) is skipped for unchanged files. This replaces the full ``rglob`` re-read on rebuilds. Excalidraw diagrams (#86, like PDFs since BUG-040) are deferred: the scan only records the title and sets ``excalidraw_text_pending``; the expensive JSON/lz-string text extraction runs in ``enrich_pdf_texts()`` after the index is queryable. Args: vault_name: Display name of the vault. vault_path: Absolute filesystem path to the vault root. vault_cfg: Optional vault configuration dict (unused for indexing, kept for compatibility). previous_files: Optional ``{relative_path: file_info}`` snapshot from a previous scan used for differential reuse. Returns: Dict with keys ``files`` (list), ``tags`` (counter dict), ``path`` (str), ``paths`` (list) and ``reused`` (int, differential hits). """ vault_root = Path(vault_path) files: list[dict[str, Any]] = [] tag_counts: dict[str, int] = {} paths: list[dict[str, str]] = [] reused = 0 if not vault_root.exists(): logger.warning(f"Vault path does not exist: {vault_path}") return {"files": [], "tags": {}, "path": vault_path, "paths": [], "reused": 0} root_resolved = vault_root.resolve(strict=False) # BUG-032: walk without following symlinks and refuse any symlink that # escapes the vault root, so external data can never be indexed/exposed. for dirpath, dirnames, filenames in os.walk(vault_root, followlinks=False): current_dir = Path(dirpath) # Prune ignored and symlinked directories in place (no recursion). dirnames[:] = [ d for d in dirnames if d not in IGNORED_DIRS and not (current_dir / d).is_symlink() ] for d in dirnames: dpath = current_dir / d rel_path_str = str(dpath.relative_to(vault_root)).replace("\\", "/") paths.append({ "path": rel_path_str, "name": d, "type": "directory" }) for fname in filenames: fpath = current_dir / fname if fpath.is_symlink(): try: target = fpath.resolve(strict=True) except OSError: continue try: target.relative_to(root_resolved) except ValueError: logger.warning(f"Skipping symlink outside vault: {fpath}") continue rel_path_str = str(fpath.relative_to(vault_root)).replace("\\", "/") ext = fpath.suffix.lower() # Also match extensionless files named like Dockerfile, Makefile basename_lower = fpath.name.lower() if ext not in SUPPORTED_EXTENSIONS and basename_lower not in ("dockerfile", "makefile", "cmakelists.txt"): continue # Add file to path index paths.append({ "path": rel_path_str, "name": fname, "type": "file" }) try: relative = fpath.relative_to(vault_root) stat = fpath.stat() modified = datetime.fromtimestamp(stat.st_mtime, tz=timezone.utc).isoformat() # #86 differential scan: reuse the previous entry when neither # size nor mtime changed — skips the disk read + parse below. if previous_files: prev = previous_files.get(rel_path_str) if ( prev is not None and prev.get("size") == stat.st_size and prev.get("modified") == modified ): file_info = {**prev, "tags": list(prev.get("tags", []))} files.append(file_info) for tag in file_info.get("tags", []): tag_counts[tag] = tag_counts.get(tag, 0) + 1 reused += 1 # The global backlink index is rebuilt on every scan, # so re-register this file's wikilinks from its # (cached) content instead of re-reading the disk. if file_info.get("extension") == ".md" and file_info.get("content"): try: _extract_wikilinks_for_backlinks( vault_name, file_info["path"], file_info.get("title", ""), file_info["content"], ) except Exception: pass continue # PDF handling — special path (binary, uses pdf_reader) tags: list[str] = [] pdf_text_pending = False excalidraw_text_pending = False if ext == ".pdf": from backend.pdf_reader import extract_pdf_metadata # BUG-040: only the (cheap) metadata is read during the # scan. Full-text extraction is deferred to a background # pass (``enrich_pdf_texts``) so a vault with many/large # PDFs no longer blocks startup and index rebuilds. pdf_meta = extract_pdf_metadata(fpath) title = pdf_meta.get("title") or fpath.stem.replace("-", " ").replace("_", " ") raw = "" content_preview = "" pdf_text_pending = True elif ext == ".excalidraw" or fpath.name.lower().endswith(".excalidraw.md"): # #86: defer the expensive JSON/lz-string text extraction # (read + decompress + element walk) to ``enrich_pdf_texts`` # so the scan stays cheap; title comes from the filename. raw = "" title = fpath.stem.replace(".excalidraw", "").replace("-", " ").replace("_", " ") content_preview = "" excalidraw_text_pending = True elif is_media(ext): # #108 — images (and future media, #109) are binary: index # name/size/mtime only and never read the bytes. ``content`` # stays empty so the TF-IDF index remains clean. raw = "" title = fpath.stem.replace("-", " ").replace("_", " ") content_preview = "" elif ext == ".xlsx": # #153 A5 — a workbook stays rendered by the viewer, but its # cell values are now indexed as text so a spreadsheet is # findable by its content (parity with _index_single_file_sync). raw = extract_xlsx_indexable(fpath) title = fpath.stem.replace("-", " ").replace("_", " ") content_preview = raw[:200].strip() else: raw = fpath.read_text(encoding="utf-8", errors="replace") title = fpath.stem.replace("-", " ").replace("_", " ") content_preview = raw[:200].strip() if ext == ".md": post = parse_markdown_file(raw) tags = _extract_tags(post) inline_tags = _extract_inline_tags(post.content) tags = list(set(tags) | set(inline_tags)) title = _extract_title(post, fpath) content_preview = post.content[:200].strip() _extract_wikilinks_for_backlinks( vault_name, str(relative).replace("\\", "/"), title, post.content ) file_info = { "path": str(relative).replace("\\", "/"), "title": title, "tags": tags, "content_preview": content_preview, "content": raw[:SEARCH_CONTENT_LIMIT], "size": stat.st_size, "modified": modified, "extension": ext, } if pdf_text_pending: file_info["pdf_text_pending"] = True if excalidraw_text_pending: file_info["excalidraw_text_pending"] = True files.append(file_info) for tag in tags: tag_counts[tag] = tag_counts.get(tag, 0) + 1 except PermissionError: logger.debug(f"Permission denied, skipping {fpath}") continue except Exception as e: logger.error(f"Error indexing {fpath}: {e}") continue logger.info( f"Vault '{vault_name}': indexed {len(files)} files " f"({reused} reused), {len(paths)} paths, {len(tag_counts)} unique tags" ) return {"files": files, "tags": tag_counts, "path": vault_path, "paths": paths, "config": {}, "reused": reused} def _read_excalidraw_indexable_text(file_path: Path) -> str: """Read an excalidraw file and return its indexable text (blocking helper). Runs inside an executor via ``enrich_pdf_texts`` so the lz-string decompression of large diagrams never blocks the event loop. """ try: raw = file_path.read_text(encoding="utf-8", errors="replace") except OSError: return "" try: return extract_excalidraw_indexable(raw) except Exception: # pragma: no cover - defensive return "" async def enrich_pdf_texts(vault_name: str | None = None) -> int: """Extract text deferred during the scan: PDFs (BUG-040) + excalidraw (#86). ``_scan_vault`` only reads PDF metadata and excalidraw filenames so a vault with many or large heavy files starts serving immediately. This coroutine runs *after* the index (and the inverted index) is ready, extracts the missing text off the event loop and updates the in-memory entry plus the incremental index hooks. Args: vault_name: Restrict the pass to a single vault; ``None`` covers every indexed vault. Returns: Number of deferred files (PDF + excalidraw) whose text extraction was attempted. """ from backend.pdf_reader import extract_pdf_text pending: list[tuple[str, dict[str, Any], Path, str]] = [] with _index_lock: for name, vault_data in index.items(): if vault_name is not None and name != vault_name: continue vault_root = Path(vault_data.get("path", "")) for file_info in vault_data.get("files", []): if file_info.get("pdf_text_pending"): pending.append((name, file_info, vault_root / file_info["path"], "pdf")) elif file_info.get("excalidraw_text_pending"): pending.append((name, file_info, vault_root / file_info["path"], "excalidraw")) if not pending: return 0 loop = asyncio.get_running_loop() enriched = 0 for name, file_info, file_path, kind in pending: try: if kind == "pdf": raw = await loop.run_in_executor(None, extract_pdf_text, file_path, 100000) else: raw = await loop.run_in_executor(None, _read_excalidraw_indexable_text, file_path) except Exception as exc: # pragma: no cover - defensive logger.warning("Deferred text enrichment failed for %s: %s", file_path, exc) raw = "" file_info["content"] = raw[:SEARCH_CONTENT_LIMIT] file_info["content_preview"] = raw[:200].strip() file_info.pop("pdf_text_pending", None) file_info.pop("excalidraw_text_pending", None) enriched += 1 if _on_index_change: try: _on_index_change("add", name, file_info["path"], file_info) except Exception as exc: # pragma: no cover - defensive logger.warning( "Index hook failed after deferred enrichment for %s: %s", file_path, exc ) logger.info("Deferred text enrichment: extracted text for %d file(s)", enriched) return enriched async def build_index(progress_callback=None) -> None: """Build the full in-memory index for all configured vaults. Runs vault scans concurrently, inserting them incrementally into the global index. Notifies progress via the provided callback. #86 differential rebuild: the previous per-vault ``{path: file_info}`` snapshots are captured before the clear and handed to ``_scan_vault`` so unchanged files (same size + mtime) are reused without disk re-reads. """ global index, vault_config vault_config.clear() vault_config.update(load_vault_config()) # Note: vault_settings are now only used for UI display preferences (hideHiddenFiles) # Indexing always includes all files regardless of settings global _index_generation with _index_lock: previous_snapshot: dict[str, dict[str, dict[str, Any]]] = { name: {f["path"]: f for f in vdata.get("files", [])} for name, vdata in index.items() } index.clear() _file_lookup.clear() path_index.clear() _backlink_index.clear() _index_generation += 1 if not vault_config: logger.warning("No vaults configured. Set VAULT_N_NAME / VAULT_N_PATH env vars.") if progress_callback: await progress_callback("complete", {"total": 0}) return if progress_callback: await progress_callback("start", {"total_vaults": len(vault_config)}) loop = asyncio.get_event_loop() async def _process_vault(name: str, config: dict[str, Any]): import functools vault_path = config["path"] scan = functools.partial( _scan_vault, name, vault_path, config, previous_snapshot.get(name) ) vault_data = await loop.run_in_executor(None, scan) vault_data["config"] = config # Build lookup entries for the new vault new_lookup_entries: dict[str, list[dict[str, str]]] = {} for f in vault_data["files"]: entry = {"vault": name, "path": f["path"]} fname = f["path"].rsplit("/", 1)[-1].lower() fpath_lower = f["path"].lower() for key in (fname, fpath_lower): if key not in new_lookup_entries: new_lookup_entries[key] = [] new_lookup_entries[key].append(entry) async_lock = _get_async_lock() async with async_lock: with _index_lock: index[name] = vault_data for key, entries in new_lookup_entries.items(): if key not in _file_lookup: _file_lookup[key] = [] _file_lookup[key].extend(entries) path_index[name] = vault_data.get("paths", []) global _index_generation _index_generation += 1 if progress_callback: await progress_callback("progress", { "vault": name, "files": len(vault_data["files"]), "tags": len(vault_data["tags"]) }) # Run vault scans concurrently tasks = [] for name, config in vault_config.items(): tasks.append(_process_vault(name, config)) if tasks: await asyncio.gather(*tasks) # Record timestamp of full index rebuild global _last_full_index_ts _last_full_index_ts = datetime.now(timezone.utc).isoformat() # Build attachment index from backend.attachment_indexer import build_attachment_index await build_attachment_index(vault_config) total_files = sum(len(v["files"]) for v in index.values()) logger.info(f"Index built: {len(index)} vaults, {total_files} total files") if progress_callback: await progress_callback("complete", {"total_vaults": len(vault_config), "total_files": total_files}) async def reload_index() -> dict[str, Any]: """Force a full re-index of all vaults and return per-vault statistics. Returns: Dict mapping vault names to their file/tag counts. """ await build_index() # BUG-040/#86: complete the deferred PDF + excalidraw extraction. await enrich_pdf_texts() # The inverted index is NOT updated by the hooks here: the rebuild above # replaces whole vault entries, so the incremental notifications are not # emitted for the files that only changed content. Without this, a manual # reindex left TF-IDF search serving a stale index (BUG-089). from backend.search import init_inverted_index init_inverted_index() stats = {} for name, data in index.items(): stats[name] = {"file_count": len(data["files"]), "tag_count": len(data["tags"])} return stats async def reload_single_vault(vault_name: str) -> dict[str, Any]: """Force a re-index of a single vault and return its statistics. Args: vault_name: Name of the vault to reindex. Returns: Dict with vault statistics (file_count, tag_count). Raises: ValueError: If vault_name is not found in configuration. """ global vault_config # Reload vault config from env vars vault_config.update(load_vault_config()) if vault_name not in vault_config: raise ValueError(f"Vault '{vault_name}' not found in configuration") config = vault_config[vault_name] # #86 differential rescan: snapshot this vault's entries before removal so # unchanged files are reused without disk re-reads. with _index_lock: _previous = {f["path"]: f for f in index.get(vault_name, {}).get("files", [])} # Remove old vault data from index structures await remove_vault_from_index(vault_name) # Re-add the vault with updated configuration import functools vault_path = config["path"] loop = asyncio.get_event_loop() scan = functools.partial(_scan_vault, vault_name, vault_path, config, _previous) vault_data = await loop.run_in_executor(None, scan) vault_data["config"] = config # Build lookup entries for the vault new_lookup_entries: dict[str, list[dict[str, str]]] = {} for f in vault_data["files"]: entry = {"vault": vault_name, "path": f["path"]} fname = f["path"].rsplit("/", 1)[-1].lower() fpath_lower = f["path"].lower() for key in (fname, fpath_lower): if key not in new_lookup_entries: new_lookup_entries[key] = [] new_lookup_entries[key].append(entry) async_lock = _get_async_lock() async with async_lock: with _index_lock: index[vault_name] = vault_data for key, entries in new_lookup_entries.items(): if key not in _file_lookup: _file_lookup[key] = [] _file_lookup[key].extend(entries) path_index[vault_name] = vault_data.get("paths", []) global _index_generation _index_generation += 1 # Rebuild attachment index for this vault only from backend.attachment_indexer import build_attachment_index await build_attachment_index({vault_name: config}) # BUG-040/#86: complete the deferred PDF + excalidraw extraction. await enrich_pdf_texts(vault_name) # Same as reload_index: the vault entry was replaced wholesale, so rebuild # the inverted index or TF-IDF search keeps serving stale postings # (BUG-089). from backend.search import init_inverted_index init_inverted_index() stats = {"file_count": len(vault_data["files"]), "tag_count": len(vault_data["tags"])} logger.info(f"Vault '{vault_name}' reindexed: {stats['file_count']} files, {stats['tag_count']} tags") return stats def get_vault_names() -> list[str]: """Return the list of all indexed vault names.""" return list(index.keys()) def get_vault_data(vault_name: str) -> dict[str, Any] | None: """Return the full index data for a vault, or ``None`` if not found.""" return index.get(vault_name) def _get_async_lock() -> asyncio.Lock: """Get or create the async lock (must be called from an event loop).""" global _async_index_lock if _async_index_lock is None: _async_index_lock = asyncio.Lock() return _async_index_lock def _index_single_file_sync(vault_name: str, vault_path: str, file_path: str, vault_cfg: dict[str, Any] | None = None) -> dict[str, Any] | None: """Synchronously read and parse a single file for indexing. All files are indexed, including hidden files (starting with '.'). Args: vault_name: Name of the vault. vault_path: Absolute path to vault root. file_path: Absolute path to the file. vault_cfg: Optional vault configuration dict (unused for indexing, kept for compatibility). Returns: File info dict or None if the file cannot be read. """ try: fpath = Path(file_path) vault_root = Path(vault_path) if not fpath.exists() or not fpath.is_file(): return None relative = fpath.relative_to(vault_root) ext = fpath.suffix.lower() basename_lower = fpath.name.lower() if ext not in SUPPORTED_EXTENSIONS and basename_lower not in ("dockerfile", "makefile", "cmakelists.txt"): return None stat = fpath.stat() modified = datetime.fromtimestamp(stat.st_mtime, tz=timezone.utc).isoformat() # PDF handling — binary, must go through pdf_reader (same as _scan_vault) tags: list[str] = [] title = fpath.stem.replace("-", " ").replace("_", " ") if ext == ".pdf": from backend.pdf_reader import extract_pdf_metadata, extract_pdf_text raw = extract_pdf_text(fpath, max_chars=SEARCH_CONTENT_LIMIT) pdf_meta = extract_pdf_metadata(fpath) title = pdf_meta.get("title") or title content_preview = raw[:200].strip() elif ext == ".excalidraw" or fpath.name.lower().endswith(".excalidraw.md"): raw = fpath.read_text(encoding="utf-8", errors="replace") raw = extract_excalidraw_indexable(raw) title = fpath.stem.replace(".excalidraw", "").replace("-", " ").replace("_", " ") content_preview = raw[:200].strip() elif is_media(ext): # #108 — binary media: metadata only, never read the bytes. raw = "" content_preview = "" elif ext == ".xlsx": # #153 A5 — index sheet names + header rows as text (see _scan_vault). raw = extract_xlsx_indexable(fpath) content_preview = raw[:200].strip() else: raw = fpath.read_text(encoding="utf-8", errors="replace") content_preview = raw[:200].strip() if ext == ".md": post = parse_markdown_file(raw) tags = _extract_tags(post) inline_tags = _extract_inline_tags(post.content) tags = list(set(tags) | set(inline_tags)) title = _extract_title(post, fpath) content_preview = post.content[:200].strip() return { "path": str(relative).replace("\\", "/"), "title": title, "tags": tags, "content_preview": content_preview, "content": raw[:SEARCH_CONTENT_LIMIT], "size": stat.st_size, "modified": modified, "extension": ext, } except PermissionError: logger.debug(f"Permission denied: {file_path}") return None except Exception as e: logger.error(f"Error parsing file {file_path}: {e}") return None def _remove_file_from_structures(vault_name: str, rel_path: str) -> dict[str, Any] | None: """Remove a file from all index structures. Returns removed file info or None. Must be called under _index_lock or _async_index_lock. """ global _index_generation vault_data = index.get(vault_name) if not vault_data: return None # Remove from files list removed = None files = vault_data["files"] for i, f in enumerate(files): if f["path"] == rel_path: removed = files.pop(i) break if not removed: return None # Update tag counts for tag in removed.get("tags", []): tc = vault_data["tags"] if tag in tc: tc[tag] -= 1 if tc[tag] <= 0: del tc[tag] # Remove from _file_lookup fname_lower = rel_path.rsplit("/", 1)[-1].lower() fpath_lower = rel_path.lower() for key in (fname_lower, fpath_lower): entries = _file_lookup.get(key, []) _file_lookup[key] = [e for e in entries if not (e["vault"] == vault_name and e["path"] == rel_path)] if not _file_lookup[key]: del _file_lookup[key] # Remove from path_index if vault_name in path_index: path_index[vault_name] = [p for p in path_index[vault_name] if p["path"] != rel_path] _index_generation += 1 # Notify inverted index for incremental update if _on_index_change: _on_index_change('remove', vault_name, rel_path, removed) # type: ignore[misc] return removed def _ensure_parent_dirs_in_path_index(vault_name: str, rel_path: str, existing: set): """Ensure all parent directories of a file path exist in path_index. For a path like ``a/b/c/file.md``, ensures ``a``, ``a/b``, and ``a/b/c`` directory entries exist in path_index. Args: vault_name: Name of the vault. rel_path: Relative path of the file (e.g. ``a/b/c/file.md``). existing: Set of paths already in path_index for this vault. """ parts = rel_path.split("/") # Build parent directory paths for i in range(1, len(parts)): dir_path = "/".join(parts[:i]) if dir_path and dir_path not in existing: existing.add(dir_path) path_index[vault_name].append({ "path": dir_path, "name": parts[i - 1], "type": "directory", }) def _add_file_to_structures(vault_name: str, file_info: dict[str, Any]): """Add a file entry to all index structures. Must be called under _index_lock or _async_index_lock. """ global _index_generation vault_data = index.get(vault_name) if not vault_data: return vault_data["files"].append(file_info) # Update tag counts for tag in file_info.get("tags", []): vault_data["tags"][tag] = vault_data["tags"].get(tag, 0) + 1 # Add to _file_lookup rel_path = file_info["path"] fname_lower = rel_path.rsplit("/", 1)[-1].lower() fpath_lower = rel_path.lower() entry = {"vault": vault_name, "path": rel_path} for key in (fname_lower, fpath_lower): if key not in _file_lookup: _file_lookup[key] = [] _file_lookup[key].append(entry) # Add to path_index — also ensures parent directories are present if vault_name in path_index: existing = {p["path"] for p in path_index[vault_name]} if rel_path not in existing: path_index[vault_name].append({ "path": rel_path, "name": rel_path.rsplit("/", 1)[-1], "type": "file", }) # Ensure all parent directories are in path_index _ensure_parent_dirs_in_path_index(vault_name, rel_path, existing) _index_generation += 1 # Notify inverted index for incremental update if _on_index_change: _on_index_change('add', vault_name, file_info["path"], file_info) # type: ignore[misc] async def update_single_file(vault_name: str, abs_file_path: str) -> dict[str, Any] | None: """Re-index a single file without full rebuild. Reads the file, removes the old entry if present, inserts the new one. Thread-safe via async lock. Args: vault_name: Name of the vault containing the file. abs_file_path: Absolute filesystem path to the file. Returns: The new file info dict, or None if file could not be indexed. """ vault_data = index.get(vault_name) if not vault_data: logger.warning(f"update_single_file: vault '{vault_name}' not in index") return None vault_path = vault_data.get("path") or vault_config.get(vault_name, {}).get("path", "") if not vault_path: return None # Get vault configuration for hidden files handling vault_cfg = vault_data.get("config") or vault_config.get(vault_name, {}) loop = asyncio.get_event_loop() file_info = await loop.run_in_executor(None, _index_single_file_sync, vault_name, vault_path, abs_file_path, vault_cfg) lock = _get_async_lock() async with lock: # Remove old entry if exists try: rel_path = str(Path(abs_file_path).relative_to(vault_path)).replace("\\", "/") except ValueError: logger.warning(f"File {abs_file_path} not under vault {vault_path}") return None _remove_file_from_structures(vault_name, rel_path) if file_info: _add_file_to_structures(vault_name, file_info) if file_info: logger.debug(f"Updated: {vault_name}/{file_info['path']}") return file_info async def remove_single_file(vault_name: str, abs_file_path: str) -> dict[str, Any] | None: """Remove a single file from the index. Args: vault_name: Name of the vault. abs_file_path: Absolute path to the deleted file. Returns: The removed file info dict, or None if not found. """ vault_data = index.get(vault_name) if not vault_data: return None vault_path = vault_data.get("path") or vault_config.get(vault_name, {}).get("path", "") if not vault_path: return None try: rel_path = str(Path(abs_file_path).relative_to(vault_path)).replace("\\", "/") except ValueError: return None lock = _get_async_lock() async with lock: removed = _remove_file_from_structures(vault_name, rel_path) if removed: logger.debug(f"Removed: {vault_name}/{rel_path}") return removed async def handle_file_move(vault_name: str, src_abs: str, dest_abs: str) -> dict[str, Any] | None: """Handle a file move/rename by removing old entry and indexing new location. Args: vault_name: Name of the vault. src_abs: Absolute path of the source (old location). dest_abs: Absolute path of the destination (new location). Returns: The new file info dict, or None. """ await remove_single_file(vault_name, src_abs) return await update_single_file(vault_name, dest_abs) async def remove_vault_from_index(vault_name: str): """Remove an entire vault from the index. Args: vault_name: Name of the vault to remove. """ global _index_generation lock = _get_async_lock() async with lock: vault_data = index.pop(vault_name, None) if not vault_data: return # Clean _file_lookup for f in vault_data.get("files", []): rel_path = f["path"] fname_lower = rel_path.rsplit("/", 1)[-1].lower() fpath_lower = rel_path.lower() for key in (fname_lower, fpath_lower): entries = _file_lookup.get(key, []) _file_lookup[key] = [e for e in entries if e["vault"] != vault_name] if not _file_lookup[key]: _file_lookup.pop(key, None) # Notify the inverted index, otherwise every document of the vault # stays in it as a ghost (postings, doc_info, doc_vault, vault_docs) # and keeps matching searches for a vault that no longer exists. if _on_index_change: _on_index_change('remove', vault_name, rel_path, f) # type: ignore[misc] # Clean path_index path_index.pop(vault_name, None) # Clean vault_config vault_config.pop(vault_name, None) _index_generation += 1 logger.info(f"Removed vault '{vault_name}' from index") async def add_vault_to_index(vault_name: str, vault_path: str) -> dict[str, Any]: """Add a new vault to the index dynamically. Args: vault_name: Display name for the vault. vault_path: Absolute filesystem path to the vault. Returns: Dict with vault stats (file_count, tag_count). """ global _index_generation vault_config[vault_name] = { "path": vault_path, "attachmentsPath": None, "scanAttachmentsOnStartup": True, } loop = asyncio.get_event_loop() vault_data = await loop.run_in_executor(None, _scan_vault, vault_name, vault_path, vault_config[vault_name]) vault_data["config"] = vault_config[vault_name] # Build lookup entries for the new vault new_lookup_entries: dict[str, list[dict[str, str]]] = {} for f in vault_data["files"]: entry = {"vault": vault_name, "path": f["path"]} fname = f["path"].rsplit("/", 1)[-1].lower() fpath_lower = f["path"].lower() for key in (fname, fpath_lower): if key not in new_lookup_entries: new_lookup_entries[key] = [] new_lookup_entries[key].append(entry) lock = _get_async_lock() async with lock: index[vault_name] = vault_data for key, entries in new_lookup_entries.items(): if key not in _file_lookup: _file_lookup[key] = [] _file_lookup[key].extend(entries) path_index[vault_name] = vault_data.get("paths", []) _index_generation += 1 stats = {"file_count": len(vault_data["files"]), "tag_count": len(vault_data["tags"])} logger.info(f"Added vault '{vault_name}': {stats['file_count']} files, {stats['tag_count']} tags") return stats def find_file_in_index(link_target: str, current_vault: str) -> dict[str, str] | None: """Find a file matching a wikilink target using O(1) lookup table. Searches by filename first, then by full relative path. Prefers results from *current_vault* when multiple matches exist. Args: link_target: The wikilink target (e.g. ``"My Note"`` or ``"folder/My Note"``). current_vault: Name of the vault the link originates from. Returns: Dict with ``vault`` and ``path`` keys, or ``None`` if not found. """ target_lower = link_target.lower().strip() if not target_lower.endswith(".md"): target_lower += ".md" candidates = _file_lookup.get(target_lower, []) if not candidates: return None # Prefer current vault when multiple vaults contain a match for c in candidates: if c["vault"] == current_vault: return c return candidates[0] # --------------------------------------------------------------------------- # Backlink index: tracks which files link to which targets # --------------------------------------------------------------------------- def _extract_wikilinks_for_backlinks( vault_name: str, source_path: str, source_title: str, content: str ): """Extract wikilinks from markdown content and populate the backlink index. For each `[[target]]` or `[[target|display]]` found in the content, adds the source file to the backlink index of the target. Args: vault_name: The vault containing the source file. source_path: Relative path of the source file in the vault. source_title: Title of the source file. content: The markdown content to scan for wikilinks. """ global _backlink_index wikilink_pattern = re.compile(r'\[\[([^\]|#]+)(?:[|#][^\]]+)?\]\]') targets = set() for match in wikilink_pattern.finditer(content): target = match.group(1).strip() target_lower = target.lower() if not target_lower.endswith(".md"): target_lower += ".md" targets.add(target_lower) if vault_name not in _backlink_index: _backlink_index[vault_name] = {} for target in targets: if target not in _backlink_index[vault_name]: _backlink_index[vault_name][target] = [] # Avoid duplicates for same source-target pair existing = [e for e in _backlink_index[vault_name][target] if e["path"] == source_path] if not existing: _backlink_index[vault_name][target].append({ "vault": vault_name, "path": source_path, "title": source_title, }) def get_backlinks(vault_name: str, file_path: str) -> list[dict[str, str]]: """Get all files that link to the given file via wikilinks. Searches across all vaults for backlinks pointing to the target file. Args: vault_name: The vault containing the target file. file_path: Relative path of the target file in the vault. Returns: List of ``{vault, path, title}`` dicts for files linking to the target. """ global _backlink_index target_key = file_path.lower() if not target_key.endswith(".md"): target_key += ".md" results = [] for vindex in _backlink_index.values(): bl = vindex.get(target_key, []) results.extend(bl) return results def get_conflicts() -> list: """Scan all vaults for Syncthing/Nextcloud sync-conflict files. Returns: List of conflict dicts with vault, conflict_path, original_path, conflict_date, and conflict_title. """ import re conflicts = [] pattern = re.compile(r'\.sync-conflict-(\d{8}-\d{6})\.') for vname, vdata in index.items(): for f in vdata.get("files", []): m = pattern.search(f["path"]) if m: orig_path = pattern.sub("", f["path"]) conflicts.append({ "vault": vname, "conflict_path": f["path"], "original_path": orig_path, "conflict_date": m.group(1), "conflict_title": f.get("title", ""), "conflict_size": f.get("size", 0), }) return conflicts