diff --git a/CHANGELOG.md b/CHANGELOG.md index 923fd2a..6f91e45 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,17 @@ ## [Unreleased] +- **Bugfix-Sweep (2026-10-06)**: + - **FTS5-Delete-Trigger repariert (Schema v5)**: Die Trigger nutzten den FTS5-`'delete'`-Befehl auf regulären FTS5-Tabellen → `SQL logic error` bei jedem Löschen von Chunks (Re-Ingest, `index --force`, Web-Viewer-Löschen, `deduplicate`). Bestehende Datenbanken werden beim Öffnen migriert. + - **Löschen konsistent**: Web-Viewer und `deduplicate` entfernen zuerst die DB-Einträge (Transaktion) und verschieben danach die Datei in `_Papierkorb`; Fehler pro Eintrag brechen den Lauf nicht mehr ab. + - **GUI**: Falsche relative Imports (`..schema`/`..config`) ließen die Dokumentliste immer leer; nach dem Sortieren wurde das falsche Dokument geöffnet; Scan-Threads nutzten eine im GUI-Thread erzeugte SQLite-Connection (jede Datei schlug fehl); EventBus-Handler liefen in Worker-Threads (jetzt Queued-Dispatch in den GUI-Thread). + - **Summarizer**: Fehlgeschlagene Chunks führen zu Status `error` statt `done`; verwaiste `processing`-Einträge (> 30 min) werden zurückgesetzt, Ctrl+C setzt das aktuelle Item auf `pending`; Kostenschätzung auf claude-haiku-4-5-Preise ($1/$5 pro MTok) aktualisiert. + - **Suche**: FTS5-Syntaxfehler (z.B. `COVID-19`, `E-Mail`, `c++`) werden mit quotierten Tokens wiederholt, danach LIKE-Fallback (auch für Dokumente); `%`/`_` werden in LIKE-Fallbacks und Verzeichnisfiltern escaped; Verzeichnisfilter matchen keine Geschwister-Ordner mehr (`docs` ≠ `docs2`). + - **Ingest**: Re-Ingest setzt Zusammenfassungen und Queue-Status zurück; archivierte Dokumente werden über `archived_path` gefunden; UTF-8-BOM wird entfernt (Frontmatter-Erkennung); der Chunker erzwingt die Obergrenze von 500 Wörtern auch bei Text ohne Satzgrenzen. + - **Sonstiges**: `is_active = 0` bleibt beim Skill-Index erhalten; Operator-Präzedenz bei `content_hash` korrigiert; `Config` teilt Default-Listen nicht mehr zwischen Instanzen; `zoll_station.py` ruft die CLI als Modul auf; plattformübergreifendes Öffnen von Dateien (Windows/macOS/Linux). + - **Transit-Sync optional**: `sqlite-transit-sync` ist nicht auf PyPI – Import ist jetzt optional, die Sync-Tests werden ohne Paket übersprungen statt die gesamte Test-Collection abzubrechen. + - **Tests**: 43 neue Regressionstests (`tests/test_bugfixes_2026_10.py`) plus Lifecycle-Tests in `tests/test_transit.py`. + - **Discoverability, Visual Architecture & Level 1 SBOM (Pfad B, 2026-09-22)**: - **18-Point Bilingual Navigation Parity**: Implemented 18-point dual-anchor navigation parity (``...``) across `README.md` and `README_de.md` with reciprocal dual anchors for backward compatibility. - **Target Personas & SEO Discovery**: Added explicit persona profiles (`[PERSONA-01]` to `[PERSONA-04]`) with high-intent search queries for AI engineers, researchers, compliance officers, and automation builders. diff --git a/chunker.py b/chunker.py index 5b176b8..9962889 100644 --- a/chunker.py +++ b/chunker.py @@ -65,7 +65,7 @@ def split_frontmatter(text: str) -> Tuple[Optional[str], str]: if not text: return None, "" - text = text.strip() + text = text.strip().lstrip("\ufeff") # UTF-8-BOM wuerde '---' verdecken # Frontmatter: beginnt mit --- und endet mit --- if text.startswith("---"): @@ -116,6 +116,21 @@ def _split_sentences(text: str) -> List[str]: return segments +def _hard_split(segments: List[str], window: int) -> List[str]: + """Erzwingt die harte Obergrenze: Segmente > MAX_CHUNK_SIZE (z.B. lange + Einzeiler ohne Satzgrenzen) werden in Wortfenster der Groesse `window` + zerlegt.""" + result: List[str] = [] + for segment in segments: + words = segment.split() + if len(words) <= MAX_CHUNK_SIZE: + result.append(segment) + continue + for i in range(0, len(words), window): + result.append(" ".join(words[i:i + window])) + return result + + def chunk_text(text: str, chunk_size: int = DEFAULT_CHUNK_SIZE, overlap: int = DEFAULT_OVERLAP, separate_frontmatter: bool = True) -> List[Chunk]: @@ -151,8 +166,9 @@ def chunk_text(text: str, chunk_size: int = DEFAULT_CHUNK_SIZE, )) chunk_idx += 1 - # Body in Segmente splitten - segments = _split_sentences(body) + # Body in Segmente splitten (Segmente > MAX_CHUNK_SIZE hart teilen) + segments = _hard_split(_split_sentences(body), + max(1, min(chunk_size, MAX_CHUNK_SIZE))) if not segments: return chunks @@ -225,8 +241,10 @@ def chunk_text(text: str, chunk_size: int = DEFAULT_CHUNK_SIZE, content = "\n\n".join(current_parts) tokens = estimate_tokens(content) - # Zu kleiner letzter Chunk? An vorherigen anhaengen - if tokens < MIN_CHUNK_SIZE and len(chunks) > 0 and not chunks[-1].is_frontmatter: + # Zu kleiner letzter Chunk? An vorherigen anhaengen (sofern die + # harte Obergrenze dabei nicht ueberschritten wird) + if (tokens < MIN_CHUNK_SIZE and len(chunks) > 0 and not chunks[-1].is_frontmatter + and chunks[-1].token_count + tokens <= MAX_CHUNK_SIZE): prev = chunks[-1] merged = prev.content + "\n\n" + content chunks[-1] = Chunk( diff --git a/config.py b/config.py index 01e4b36..77f0ae8 100644 --- a/config.py +++ b/config.py @@ -16,6 +16,7 @@ __all__ = ["Config", "get_config"] +import copy import json import os from pathlib import Path @@ -66,7 +67,8 @@ class Config: """JSON-basierte Konfiguration.""" def __init__(self, config_path: Optional[Path] = None): - self._data = dict(DEFAULT_CONFIG) + # deepcopy: Listen (z.B. indexed_directories) nicht mit DEFAULT_CONFIG teilen + self._data = copy.deepcopy(DEFAULT_CONFIG) self._path = config_path if config_path and config_path.exists(): self._load(config_path) diff --git a/digest.py b/digest.py index 9a90acd..f3fbe0c 100644 --- a/digest.py +++ b/digest.py @@ -31,12 +31,29 @@ from .config import get_config, Config from .ingestor import DocumentIngestor from .summarizer import Summarizer +from .utils import dir_filter_sql, escape_like, fts5_quote, move_to_trash, resolve_document_path # Lazy imports fuer optionale Module _SkillIndexer = None _WikiIndexer = None +def _fts_fetchall(conn: sqlite3.Connection, sql: str, params: list) -> list: + """Fuehrt eine FTS5-Abfrage aus (params[0] = MATCH-Query). + + Bei FTS5-Syntaxfehlern durch Freitext (z.B. 'COVID-19', 'E-Mail', 'c++', + unbalancierte Anfuehrungszeichen) wird mit quotierten Tokens wiederholt. + Scheitert auch das, wird der Fehler weitergereicht (-> LIKE-Fallback). + """ + try: + return conn.execute(sql, params).fetchall() + except sqlite3.OperationalError: + safe = fts5_quote(params[0]) + if not safe or safe == params[0]: + raise + return conn.execute(sql, [safe, *params[1:]]).fetchall() + + def _get_skill_indexer(): global _SkillIndexer if _SkillIndexer is None: @@ -190,7 +207,7 @@ def search(self, query: str, *, limit: int = 20, LIMIT ? """ - rows = conn.execute(sql, params).fetchall() + rows = _fts_fetchall(conn, sql, params) results = [] seen_skills = set() @@ -225,8 +242,8 @@ def _search_fallback(self, conn: sqlite3.Connection, query: str, limit: int, skill_type: Optional[str], category: Optional[str]) -> List[Dict[str, Any]]: """LIKE-basierte Fallback-Suche wenn FTS5 fehlschlaegt.""" - where_parts = ["(sc.content LIKE ? OR si.skill_name LIKE ?)"] - like = f"%{query}%" + where_parts = ["(sc.content LIKE ? ESCAPE '\\' OR si.skill_name LIKE ? ESCAPE '\\')"] + like = f"%{escape_like(query)}%" params: list = [like, like] if skill_type: @@ -489,9 +506,10 @@ def get_directories(self) -> List[Dict[str, Any]]: conn = self._get_conn() result = [] for d in dirs: + where_sql, where_params = dir_filter_sql("source_dir", d) count = conn.execute( - "SELECT COUNT(*) FROM documents WHERE source_dir LIKE ?", - (d + "%",) + f"SELECT COUNT(*) FROM documents WHERE {where_sql}", + where_params ).fetchone()[0] result.append({"path": d, "doc_count": count}) conn.close() @@ -611,7 +629,7 @@ def _search_documents(self, query: str, *, """FTS5-Suche ueber ingested Dokumente.""" conn = self._get_conn() try: - rows = conn.execute(""" + rows = _fts_fetchall(conn, """ SELECT d.filename, d.file_type, @@ -625,7 +643,7 @@ def _search_documents(self, query: str, *, WHERE document_fts MATCH ? ORDER BY document_fts.rank LIMIT ? - """, (query, limit)).fetchall() + """, [query, limit]) results = [] seen = set() @@ -644,10 +662,36 @@ def _search_documents(self, query: str, *, }) return results except Exception: - return [] + # FTS5 Fallback: LIKE-Suche + return self._search_documents_fallback(conn, query, limit) finally: conn.close() + def _search_documents_fallback(self, conn: sqlite3.Connection, query: str, + limit: int) -> List[Dict[str, Any]]: + """LIKE-basierte Fallback-Suche ueber Dokumente wenn FTS5 fehlschlaegt.""" + if not query.strip(): + return [] + like = f"%{escape_like(query)}%" + try: + rows = conn.execute(""" + SELECT DISTINCT d.filename, d.file_type, d.word_count + FROM documents d + LEFT JOIN document_chunks dc ON dc.doc_id = d.id + WHERE dc.content LIKE ? ESCAPE '\\' OR d.filename LIKE ? ESCAPE '\\' + LIMIT ? + """, (like, like, limit)).fetchall() + return [{ + 'source': 'document', + 'name': r['filename'], + 'type': r['file_type'], + 'snippet': '(LIKE-Fallback)', + 'relevance': 0, + 'word_count': r['word_count'], + } for r in rows] + except sqlite3.Error: + return [] + # ================================================================== # WIKI-ARTIKEL INDEXIERUNG # ================================================================== @@ -743,7 +787,7 @@ def search_wikis(self, query: str, *, limit: int = 20, LIMIT ? """ - rows = conn.execute(sql, params).fetchall() + rows = _fts_fetchall(conn, sql, params) results = [] seen_wikis = set() @@ -777,8 +821,8 @@ def search_wikis(self, query: str, *, limit: int = 20, def _search_wikis_fallback(self, conn: sqlite3.Connection, query: str, limit: int, category: Optional[str]) -> List[Dict[str, Any]]: """LIKE-basierte Fallback-Suche wenn FTS5 fehlschlaegt.""" - where_parts = ["(wc.content LIKE ? OR wi.title LIKE ?)"] - like = f"%{query}%" + where_parts = ["(wc.content LIKE ? ESCAPE '\\' OR wi.title LIKE ? ESCAPE '\\')"] + like = f"%{escape_like(query)}%" params: list = [like, like] if category: @@ -1305,14 +1349,13 @@ def main(): print("Keine Duplikate gefunden!") else: print(f"{len(dups)} Gruppen von Duplikaten gefunden.\n") - import shutil import os trash_dir = kd.db_path.parent / "_Papierkorb" trash_dir.mkdir(parents=True, exist_ok=True) for row in dups: hash_val = row['content_hash'] - docs = conn.execute("SELECT id, file_path, filename FROM documents WHERE content_hash=? ORDER BY id", (hash_val,)).fetchall() + docs = conn.execute("SELECT id, file_path, filename, archived_path FROM documents WHERE content_hash=? ORDER BY id", (hash_val,)).fetchall() print(f"\n--- Duplikat-Gruppe ({len(docs)} Dateien) ---") for i, d in enumerate(docs): print(f" {i+1}: {d['file_path']}") @@ -1324,28 +1367,28 @@ def main(): for d in docs[1:]: doc_id = d['id'] - file_path = d['file_path'] + file_path = resolve_document_path(d) print(f"-> Loesche: {file_path}") - - if os.path.exists(file_path): - base_name = os.path.basename(file_path) - trash_path = trash_dir / base_name - counter = 1 - while trash_path.exists(): - name, ext = os.path.splitext(base_name) - trash_path = trash_dir / f"{name}_{counter}{ext}" - counter += 1 + + # Erst DB-Eintrag in einer Transaktion entfernen, dann + # die Datei verschieben: schlaegt das Verschieben fehl, + # bleibt die DB trotzdem konsistent. + try: + with conn: + conn.execute("DELETE FROM summaries WHERE source_type='document' AND source_id=?", (doc_id,)) + conn.execute("DELETE FROM document_keywords WHERE doc_id=?", (doc_id,)) + conn.execute("DELETE FROM document_chunks WHERE doc_id=?", (doc_id,)) + conn.execute("DELETE FROM digest_queue WHERE source_type='document' AND source_id=?", (doc_id,)) + conn.execute("DELETE FROM documents WHERE id=?", (doc_id,)) + except sqlite3.Error as e: + print(f" Fehler beim Loeschen aus der DB (Datei unveraendert): {e}") + continue + + if file_path and os.path.exists(file_path): try: - shutil.move(file_path, trash_path) + move_to_trash(file_path, trash_dir) except Exception as e: - print(f" Fehler beim Verschieben: {e}") - - conn.execute("DELETE FROM summaries WHERE source_type='document' AND source_id=?", (doc_id,)) - conn.execute("DELETE FROM document_keywords WHERE doc_id=?", (doc_id,)) - conn.execute("DELETE FROM document_chunks WHERE doc_id=?", (doc_id,)) - conn.execute("DELETE FROM digest_queue WHERE source_type='document' AND source_id=?", (doc_id,)) - conn.execute("DELETE FROM documents WHERE id=?", (doc_id,)) - conn.commit() + print(f" Fehler beim Verschieben (DB-Eintrag bereits entfernt): {e}") print("\nDuplikat-Bereinigung abgeschlossen!") finally: conn.close() diff --git a/extractor.py b/extractor.py index 21e0dc5..9f60e90 100644 --- a/extractor.py +++ b/extractor.py @@ -170,8 +170,8 @@ def _extract_text(self, path: Path) -> ExtractedText: @staticmethod def _read_file(path: Path) -> Optional[str]: - """Liest Datei mit UTF-8, Fallback auf cp1252 (Windows).""" - for encoding in ('utf-8', 'cp1252', 'latin-1'): + """Liest Datei mit UTF-8 (BOM wird entfernt), Fallback auf cp1252 (Windows).""" + for encoding in ('utf-8-sig', 'cp1252', 'latin-1'): try: return path.read_text(encoding=encoding) except (UnicodeDecodeError, UnicodeError): diff --git a/gemini_flash_summarizer.py b/gemini_flash_summarizer.py index f6bbc8b..f16982d 100644 --- a/gemini_flash_summarizer.py +++ b/gemini_flash_summarizer.py @@ -23,9 +23,9 @@ from typing import Dict, List, Optional, Any try: - from .schema import ensure_schema + from .schema import ThreadLocalConnections, reset_stale_processing except ImportError: # script mode (python digest.py ...) - from schema import ensure_schema + from schema import ThreadLocalConnections, reset_stale_processing _DEFAULT_MODEL = "gemini-2.0-flash" @@ -59,15 +59,13 @@ class GeminiFlashSummarizer: def __init__(self, knowledge_db: Path, api_key: Optional[str] = None, model: Optional[str] = None): self.knowledge_db = knowledge_db - self._conn: Optional[sqlite3.Connection] = None + self._conns = ThreadLocalConnections() # eine Connection pro Thread self._api_key = api_key or os.environ.get('GEMINI_API_KEY') self.model = model or os.environ.get('GEMINI_MODEL') or _DEFAULT_MODEL self._client = None def _get_conn(self) -> sqlite3.Connection: - if self._conn is None: - self._conn = ensure_schema(self.knowledge_db) - return self._conn + return self._conns.get(self.knowledge_db) def _get_client(self): if self._client is None: @@ -81,9 +79,7 @@ def _get_client(self): return self._client def close(self): - if self._conn: - self._conn.close() - self._conn = None + self._conns.close_all() def summarize_queue(self, *, limit: int = 10, delay: float = 0.5) -> Dict[str, Any]: start = time.time() @@ -97,6 +93,9 @@ def summarize_queue(self, *, limit: int = 10, delay: float = 0.5) -> Dict[str, A 'items': [], } + # Nach Absturz/Abbruch haengengebliebene Items wieder freigeben + stats['reset_stale'] = reset_stale_processing(conn) + queue_items = conn.execute(""" SELECT id, source_type, source_id FROM digest_queue @@ -136,10 +135,12 @@ def summarize_queue(self, *, limit: int = 10, delay: float = 0.5) -> Dict[str, A 'output_tokens': 0, } + chunk_errors = [] for chunk_index, chunk_content in chunks: summary_result = self._summarize_chunk(chunk_content) if summary_result.get('error'): + chunk_errors.append((chunk_index, summary_result['error'])) continue conn.execute(""" @@ -168,13 +169,37 @@ def summarize_queue(self, *, limit: int = 10, delay: float = 0.5) -> Dict[str, A if delay > 0: time.sleep(delay) + stats['total_input_tokens'] += item_result['input_tokens'] + stats['total_output_tokens'] += item_result['output_tokens'] + stats['items'].append(item_result) + + if chunk_errors: + # Mindestens ein Chunk fehlgeschlagen -> 'error' statt + # 'done' (erfolgreiche Summaries bleiben erhalten) + first_idx, first_err = chunk_errors[0] + item_result['chunk_errors'] = len(chunk_errors) + conn.execute( + "UPDATE digest_queue SET status='error', error_msg=?, finished_at=CURRENT_TIMESTAMP WHERE id=?", + (f"{len(chunk_errors)}/{len(chunks)} Chunks fehlgeschlagen " + f"(Chunk {first_idx}: {first_err})"[:500], queue_id) + ) + conn.commit() + stats['errors'] += 1 + continue + conn.execute("UPDATE digest_queue SET status='done', finished_at=CURRENT_TIMESTAMP WHERE id=?", (queue_id,)) conn.commit() stats['processed'] += 1 - stats['total_input_tokens'] += item_result['input_tokens'] - stats['total_output_tokens'] += item_result['output_tokens'] - stats['items'].append(item_result) + + except KeyboardInterrupt: + # Abbruch (Ctrl+C): aktuelles Item wieder freigeben + conn.execute( + "UPDATE digest_queue SET status='pending', started_at=NULL WHERE id=?", + (queue_id,) + ) + conn.commit() + raise except Exception as e: conn.execute( diff --git a/gui/event_bus.py b/gui/event_bus.py index fdd4b40..cfd4671 100644 --- a/gui/event_bus.py +++ b/gui/event_bus.py @@ -5,7 +5,7 @@ Adaptiert von LitZentrum (src/core/event_bus.py). """ -from PySide6.QtCore import QObject, Signal +from PySide6.QtCore import QObject, QThread, Qt, Signal from enum import Enum from typing import Any, Callable, Dict, List @@ -43,12 +43,16 @@ class EventBus(QObject): """Zentrale Event-Verteilung (Singleton)""" event_fired = Signal(str, object) + # Intern: Events aus Worker-Threads in den GUI-Thread weiterreichen + _dispatch_requested = Signal(object, object) _instance = None def __init__(self): super().__init__() self._handlers: Dict[EventType, List[Callable]] = {} + self._dispatch_requested.connect( + self._dispatch, Qt.ConnectionType.QueuedConnection) @classmethod def instance(cls) -> "EventBus": @@ -68,9 +72,19 @@ def unsubscribe(self, event_type: EventType, handler: Callable): self._handlers[event_type].remove(handler) def emit(self, event_type: EventType, data: Any = None): + """Verteilt ein Event an alle Handler -- immer im Thread des Bus + (GUI-Thread). Aus Worker-Threads (z.B. Scan) wird das Event per + Queued-Signal in den GUI-Thread verschoben, damit Handler dort + gefahrlos Qt-Widgets anfassen koennen.""" + if QThread.currentThread() != self.thread(): + self._dispatch_requested.emit(event_type, data) + return + self._dispatch(event_type, data) + + def _dispatch(self, event_type: EventType, data: Any = None): self.event_fired.emit(event_type.value, data) if event_type in self._handlers: - for handler in self._handlers[event_type]: + for handler in list(self._handlers[event_type]): try: handler(data) except Exception as e: diff --git a/gui/panels/document_list.py b/gui/panels/document_list.py index bdf7ccc..b4628a4 100644 --- a/gui/panels/document_list.py +++ b/gui/panels/document_list.py @@ -12,6 +12,7 @@ from PySide6.QtCore import Qt from ..event_bus import EventType, get_event_bus +from ...utils import dir_filter_sql, open_path, resolve_document_path class DocumentListPanel(QWidget): @@ -80,7 +81,7 @@ def show_search_results(self, results): self._table.insertRow(row) name = r.get("filename", r.get("name", "?")) - self._table.setItem(row, 0, QTableWidgetItem(name)) + self._table.setItem(row, 0, self._name_item(name, len(self._docs) - 1)) self._table.setItem(row, 1, QTableWidgetItem(r.get("file_type", r.get("source", "")))) wc = QTableWidgetItem() @@ -105,23 +106,24 @@ def _load_documents(self, dir_filter=None): self._docs = [] try: - from ..schema import ensure_schema + from ...schema import ensure_schema conn = ensure_schema(self._kd.db_path) if dir_filter: - rows = conn.execute(""" + where_sql, where_params = dir_filter_sql("d.source_dir", dir_filter) + rows = conn.execute(f""" SELECT d.id, d.filename, d.file_type, d.word_count, d.chunk_count, - d.file_path, d.source_dir, + d.file_path, d.archived_path, d.source_dir, (SELECT COUNT(*) FROM summaries s WHERE s.source_type='document' AND s.source_id=d.id) as sum_count FROM documents d - WHERE d.source_dir LIKE ? + WHERE {where_sql} ORDER BY d.filename - """, (dir_filter + "%",)).fetchall() + """, where_params).fetchall() else: rows = conn.execute(""" SELECT d.id, d.filename, d.file_type, d.word_count, d.chunk_count, - d.file_path, d.source_dir, + d.file_path, d.archived_path, d.source_dir, (SELECT COUNT(*) FROM summaries s WHERE s.source_type='document' AND s.source_id=d.id) as sum_count FROM documents d @@ -136,7 +138,7 @@ def _load_documents(self, dir_filter=None): r = self._table.rowCount() self._table.insertRow(r) - self._table.setItem(r, 0, QTableWidgetItem(doc["filename"])) + self._table.setItem(r, 0, self._name_item(doc["filename"], len(self._docs) - 1)) self._table.setItem(r, 1, QTableWidgetItem(doc.get("file_type", ""))) wc = QTableWidgetItem() @@ -156,19 +158,36 @@ def _load_documents(self, dir_filter=None): self._table.setSortingEnabled(True) + @staticmethod + def _name_item(text, doc_index): + """Spalte-0-Item; merkt sich den Index in self._docs (UserRole), + da die sichtbare Zeile nach dem Sortieren nicht mehr dem Index entspricht.""" + item = QTableWidgetItem(text) + item.setData(Qt.ItemDataRole.UserRole, doc_index) + return item + + def _doc_at(self, row): + """Dokument zur sichtbaren Tabellenzeile (sortierfest).""" + item = self._table.item(row, 0) + if item is None: + return None + idx = item.data(Qt.ItemDataRole.UserRole) + if isinstance(idx, int) and 0 <= idx < len(self._docs): + return self._docs[idx] + return None + def _on_selection_changed(self, row, col, prev_row, prev_col): - if 0 <= row < len(self._docs): - doc = self._docs[row] + doc = self._doc_at(row) + if doc is not None: self._bus.emit(EventType.DOCUMENT_SELECTED, doc) def _on_double_click(self, row, col): """Doppelklick oeffnet Datei.""" - if 0 <= row < len(self._docs): - doc = self._docs[row] - file_path = doc.get("file_path", "") + doc = self._doc_at(row) + if doc is not None: + file_path = resolve_document_path(doc) if file_path: - import os try: - os.startfile(file_path) + open_path(file_path) except Exception: pass diff --git a/gui/panels/preview_panel.py b/gui/panels/preview_panel.py index 7c3c01d..f24df9b 100644 --- a/gui/panels/preview_panel.py +++ b/gui/panels/preview_panel.py @@ -5,9 +5,6 @@ Adaptiert von DokuZentrum (gui/panels/preview_panel.py). """ -import os -import sys -import subprocess from pathlib import Path from typing import Optional @@ -18,6 +15,8 @@ from PySide6.QtCore import Qt, QSize from PySide6.QtGui import QPixmap, QImage +from ...utils import open_path, resolve_document_path + class PreviewPanel(QWidget): """Rechtes Panel: Dokumentvorschau.""" @@ -116,7 +115,8 @@ def _setup_ui(self): def show_document(self, doc_data): """Zeigt Vorschau fuer ein Dokument (dict aus DB).""" self._current_doc = doc_data - file_path = doc_data.get("file_path", "") + # Archivierte Dokumente liegen unter archived_path statt file_path + file_path = resolve_document_path(doc_data) or "" self._current_path = file_path name = doc_data.get("filename", Path(file_path).name if file_path else "?") @@ -244,7 +244,7 @@ def _show_info(self, doc): doc_id = doc.get("id") if doc_id: try: - from ..schema import ensure_schema + from ...schema import ensure_schema conn = ensure_schema(self._current_doc_db_path()) # Keywords kws = conn.execute( @@ -284,7 +284,7 @@ def _show_chunks(self, doc): self._chunks_widget.setPlainText("Keine Chunks verfuegbar.") return try: - from ..schema import ensure_schema + from ...schema import ensure_schema conn = ensure_schema(self._current_doc_db_path()) chunks = conn.execute( "SELECT chunk_index, content FROM document_chunks WHERE doc_id=? ORDER BY chunk_index", @@ -304,14 +304,9 @@ def _show_chunks(self, doc): def _current_doc_db_path(self): """Gibt den DB-Pfad zurueck.""" - from ..config import get_config + from ...config import get_config return get_config().get_db_path() def _open_external(self): if self._current_path: - if sys.platform == "win32": - os.startfile(self._current_path) - elif sys.platform == "darwin": - subprocess.run(["open", self._current_path]) - else: - subprocess.run(["xdg-open", self._current_path]) + open_path(self._current_path) diff --git a/indexer.py b/indexer.py index 6bcf2cc..63c5e3c 100644 --- a/indexer.py +++ b/indexer.py @@ -18,15 +18,15 @@ import sqlite3 from pathlib import Path -from typing import Dict, Optional, Any, Union +from typing import Dict, Any, Union # Relative Imports (Paket-Kontext) mit Fallback auf absolute Imports (sys.path) try: - from .schema import ensure_schema + from .schema import ThreadLocalConnections from .chunker import chunk_text, estimate_tokens from .utils import sha256_hash as _sha256, extract_keywords as _extract_keywords except ImportError: - from schema import ensure_schema # type: ignore[no-redef] + from schema import ThreadLocalConnections # type: ignore[no-redef] from chunker import chunk_text, estimate_tokens # type: ignore[no-redef] from utils import sha256_hash as _sha256, extract_keywords as _extract_keywords # type: ignore[no-redef] @@ -42,19 +42,15 @@ class SkillIndexer: def __init__(self, knowledge_db: Path): self.knowledge_db = knowledge_db - self._conn: Optional[sqlite3.Connection] = None + self._conns = ThreadLocalConnections() # eine Connection pro Thread def _get_conn(self) -> sqlite3.Connection: """Lazy-init der DB-Connection mit Schema-Sicherstellung.""" - if self._conn is None: - self._conn = ensure_schema(self.knowledge_db) - return self._conn + return self._conns.get(self.knowledge_db) def close(self): """Schliesst DB-Connection.""" - if self._conn: - self._conn.close() - self._conn = None + self._conns.close_all() def index_from_bach(self, bach_db_path: Path, *, chunk_size: int = 350, @@ -119,7 +115,7 @@ def index_from_bach(self, bach_db_path: Path, *, for row in rows: skill_name = row['name'] content = row['content'] or '' - content_hash = row['content_hash'] or _sha256(content) if content else '' + content_hash = row['content_hash'] or (_sha256(content) if content else '') # Skip wenn kein Content if not content.strip(): @@ -162,7 +158,8 @@ def index_from_bach(self, bach_db_path: Path, *, """, ( skill_name, row['type'], row['category'], row['path'], row['description'], row['version'], content_hash, - word_count, chunk_count, row['is_active'] or 1, + word_count, chunk_count, + 1 if row['is_active'] is None else row['is_active'], str(bach_db_path), )) diff --git a/ingestor.py b/ingestor.py index 73d0bef..8a01e2d 100644 --- a/ingestor.py +++ b/ingestor.py @@ -27,7 +27,7 @@ from datetime import datetime from typing import Dict, Optional, Any -from .schema import ensure_schema +from .schema import ThreadLocalConnections from .chunker import chunk_text, estimate_tokens from .extractor import TextExtractor from .utils import sha256_hash, extract_keywords @@ -48,7 +48,7 @@ class DocumentIngestor: def __init__(self, knowledge_db: Path, config=None): self.knowledge_db = knowledge_db - self._conn: Optional[sqlite3.Connection] = None + self._conns = ThreadLocalConnections() # eine Connection pro Thread self._extractor = TextExtractor() if config: self.inbox_dir = config.get_inbox_dir() @@ -59,15 +59,11 @@ def __init__(self, knowledge_db: Path, config=None): def _get_conn(self) -> sqlite3.Connection: """Lazy-init der DB-Connection mit Schema-Sicherstellung.""" - if self._conn is None: - self._conn = ensure_schema(self.knowledge_db) - return self._conn + return self._conns.get(self.knowledge_db) def close(self): """Schliesst DB-Connection.""" - if self._conn: - self._conn.close() - self._conn = None + self._conns.close_all() def _ensure_dirs(self): """Erstellt inbox/ und archive/ Ordner falls noetig.""" @@ -182,6 +178,7 @@ def ingest_file(self, path: Path, *, page_count=excluded.page_count, extraction_method=excluded.extraction_method, source_dir=excluded.source_dir, + archived_path=NULL, ingested_at=CURRENT_TIMESTAMP """, ( str(path.resolve()), @@ -230,11 +227,24 @@ def ingest_file(self, path: Path, *, ) kw_count += 1 - # Queue-Eintrag fuer Summarization + # Re-Ingestion eines geaenderten Dokuments: alte Summaries gehoeren + # zum alten Inhalt -> loeschen und Queue-Eintrag zuruecksetzen + conn.execute( + "DELETE FROM summaries WHERE source_type='document' AND source_id=?", + (doc_id,) + ) + + # Queue-Eintrag fuer Summarization (Upsert: bestehender Eintrag + # wird wieder auf 'pending' gesetzt) conn.execute(""" - INSERT OR IGNORE INTO digest_queue + INSERT INTO digest_queue (source_type, source_id, status, step) VALUES ('document', ?, 'pending', 'summarize') + ON CONFLICT(source_type, source_id, step) DO UPDATE SET + status='pending', + started_at=NULL, + finished_at=NULL, + error_msg=NULL """, (doc_id,)) conn.commit() diff --git a/pyproject.toml b/pyproject.toml index d6ccdfb..499f3fe 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,5 +1,5 @@ [build-system] -requires = ["setuptools>=69", "wheel"] +requires = ["setuptools>=77", "wheel"] build-backend = "setuptools.build_meta" [project] @@ -8,7 +8,7 @@ version = "0.4.0" description = "Portable local knowledge database with FTS5 search, chunking, desktop GUI, and web viewer." readme = "README.md" requires-python = ">=3.10" -license = { text = "MIT" } +license = "MIT" license-files = ["LICENSE", "NOTICE", "THIRD_PARTY_LICENSES.md"] authors = [ { name = "Lukas Geiger" } @@ -43,7 +43,6 @@ classifiers = [ "Programming Language :: Python :: 3.11", "Programming Language :: Python :: 3.12", "Programming Language :: Python :: 3.13", - "License :: OSI Approved :: MIT License", "Operating System :: OS Independent", "Operating System :: Microsoft :: Windows", "Operating System :: POSIX :: Linux", diff --git a/requirements.txt b/requirements.txt index fc19352..69dfa7b 100644 --- a/requirements.txt +++ b/requirements.txt @@ -16,3 +16,6 @@ anthropic>=0.79.0 # Gemini Flash Summarization (optional, nur für summarize --flash) google-genai>=1.0.0 + +# Cross-System Transit Sync (optional, nur für transit.py / create_knowledge_sync) +# Nicht auf PyPI: pip install "git+https://github.com/ellmos-ai/sqlite-transit-sync" diff --git a/schema.py b/schema.py index 5067ac3..1350c26 100644 --- a/schema.py +++ b/schema.py @@ -16,12 +16,14 @@ weil FTS5 ueber einzelne Chunks performanter ist als ueber JSON-Blobs. """ -__all__ = ["SCHEMA_SQL", "SCHEMA_VERSION", "ensure_schema", "get_schema_version"] +__all__ = ["SCHEMA_SQL", "SCHEMA_VERSION", "STALE_PROCESSING_MINUTES", "ThreadLocalConnections", + "ensure_schema", "get_schema_version", "reset_stale_processing"] import sqlite3 +import threading from pathlib import Path -SCHEMA_VERSION = 4 +SCHEMA_VERSION = 5 SCHEMA_SQL = """ -- ======================================================================== @@ -71,11 +73,10 @@ new.content; END; +-- skill_fts ist eine regulaere FTS5-Tabelle (kein external content): +-- der FTS5-'delete'-Befehl ist dort nicht erlaubt -> direktes DELETE. CREATE TRIGGER IF NOT EXISTS chunk_ad AFTER DELETE ON skill_chunks BEGIN - INSERT INTO skill_fts(skill_fts, rowid, skill_name, content) - SELECT 'delete', old.id, - (SELECT skill_name FROM skill_index WHERE id = old.skill_id), - old.content; + DELETE FROM skill_fts WHERE rowid = old.id; END; -- Schluesselwoerter pro Skill @@ -134,11 +135,7 @@ END; CREATE TRIGGER IF NOT EXISTS wiki_chunk_ad AFTER DELETE ON wiki_chunks BEGIN - INSERT INTO wiki_fts(wiki_fts, rowid, wiki_path, title, content) - SELECT 'delete', old.id, - (SELECT wiki_path FROM wiki_index WHERE id = old.wiki_id), - (SELECT title FROM wiki_index WHERE id = old.wiki_id), - old.content; + DELETE FROM wiki_fts WHERE rowid = old.id; END; -- Schluesselwoerter pro Wiki @@ -198,10 +195,7 @@ END; CREATE TRIGGER IF NOT EXISTS doc_chunk_ad AFTER DELETE ON document_chunks BEGIN - INSERT INTO document_fts(document_fts, rowid, filename, content) - SELECT 'delete', old.id, - (SELECT filename FROM documents WHERE id = old.doc_id), - old.content; + DELETE FROM document_fts WHERE rowid = old.id; END; -- Schluesselwoerter pro Dokument @@ -324,13 +318,35 @@ """ -def ensure_schema(db_path: Path) -> sqlite3.Connection: +# Delete-Trigger, die bis Schema v4 den FTS5-'delete'-Befehl auf regulaeren +# FTS5-Tabellen nutzten (-> "SQL logic error" bei jedem DELETE von Chunks). +_FTS_DELETE_TRIGGERS = ("chunk_ad", "wiki_chunk_ad", "doc_chunk_ad") + + +def _migrate_fts_delete_triggers(conn: sqlite3.Connection) -> None: + """Entfernt veraltete FTS-Delete-Trigger (v4), damit SCHEMA_SQL sie neu anlegt. + + CREATE TRIGGER IF NOT EXISTS ersetzt bestehende Trigger nicht -- daher + werden nur die alten Varianten (erkennbar an 'delete') gedroppt. + """ + for name in _FTS_DELETE_TRIGGERS: + row = conn.execute( + "SELECT sql FROM sqlite_master WHERE type='trigger' AND name=?", + (name,) + ).fetchone() + if row and "'delete'" in (row[0] or ""): + conn.execute(f"DROP TRIGGER IF EXISTS {name}") + + +def ensure_schema(db_path: Path, *, check_same_thread: bool = True) -> sqlite3.Connection: """Erstellt Schema falls noetig, gibt Connection zurueck.""" - conn = sqlite3.connect(str(db_path), timeout=30) + conn = sqlite3.connect(str(db_path), timeout=30, + check_same_thread=check_same_thread) conn.row_factory = sqlite3.Row conn.execute("PRAGMA journal_mode=WAL") conn.execute("PRAGMA foreign_keys=ON") conn.execute("PRAGMA busy_timeout=30000") + _migrate_fts_delete_triggers(conn) conn.executescript(SCHEMA_SQL) conn.execute( "INSERT OR REPLACE INTO schema_meta (key, value) VALUES ('version', ?)", @@ -351,3 +367,70 @@ def get_schema_version(db_path: Path) -> int: return int(row[0]) if row else 0 except Exception: return 0 + + +class ThreadLocalConnections: + """Eine SQLite-Connection pro Thread. + + sqlite3-Connections duerfen nur im erzeugenden Thread benutzt werden. + Klassen, die ihre Connection cachen (Ingestor, Summarizer, Indexer), + werden aber aus GUI- und Scan-Threads aufgerufen -- daher bekommt jeder + Thread seine eigene Connection. close_all() schliesst alle; Connections + beendeter Threads werden beim naechsten get() aufgeraeumt. (Dafuer werden + die Connections mit check_same_thread=False geoeffnet, aber weiterhin nur + vom jeweiligen Thread benutzt.) + """ + + def __init__(self) -> None: + self._local = threading.local() + self._lock = threading.Lock() + self._all: list[tuple[threading.Thread, sqlite3.Connection]] = [] + + def get(self, db_path: Path) -> sqlite3.Connection: + conn = getattr(self._local, "conn", None) + if conn is None: + conn = ensure_schema(db_path, check_same_thread=False) + self._local.conn = conn + with self._lock: + dead = [c for t, c in self._all if not t.is_alive()] + self._all = [(t, c) for t, c in self._all if t.is_alive()] + self._all.append((threading.current_thread(), conn)) + self._close(dead) + return conn + + def close_all(self) -> None: + with self._lock: + conns = [c for _, c in self._all] + self._all = [] + self._local = threading.local() + self._close(conns) + + @staticmethod + def _close(conns: list[sqlite3.Connection]) -> None: + for conn in conns: + try: + conn.close() + except sqlite3.Error: + pass + + +# Queue-Items, die laenger als diese Zeit auf 'processing' stehen, gelten als +# verwaist (Absturz/Abbruch) und werden wieder auf 'pending' gesetzt. +STALE_PROCESSING_MINUTES = 30 + + +def reset_stale_processing(conn: sqlite3.Connection, + minutes: int = STALE_PROCESSING_MINUTES) -> int: + """Setzt verwaiste 'processing'-Eintraege der digest_queue auf 'pending'. + + Verwaist = started_at fehlt oder liegt mehr als `minutes` zurueck. + Gibt die Anzahl zurueckgesetzter Eintraege zurueck. + """ + cur = conn.execute( + "UPDATE digest_queue SET status='pending', started_at=NULL " + "WHERE status='processing' " + "AND (started_at IS NULL OR started_at < datetime('now', ?))", + (f"-{int(minutes)} minutes",) + ) + conn.commit() + return cur.rowcount diff --git a/summarizer.py b/summarizer.py index 04ceeaa..6551dac 100644 --- a/summarizer.py +++ b/summarizer.py @@ -43,7 +43,7 @@ from pathlib import Path from typing import Dict, List, Optional, Any, Callable -from .schema import ensure_schema +from .schema import ThreadLocalConnections, reset_stale_processing # System-Prompt fuer Summarization _SYSTEM_PROMPT = """\ @@ -62,6 +62,9 @@ "ollama": "qwen3:4b", } +# Kosten-SCHAETZUNG Claude Haiku 4.5 in USD pro Million Tokens (input, output) +_ANTHROPIC_PRICE_PER_MTOK = (1.00, 5.00) + class Summarizer: """LLM-basierte Chunk-Summarization (Provider-agnostisch). @@ -87,7 +90,7 @@ def __init__( self.model = model or _DEFAULT_MODELS.get(provider, "") self.base_url = base_url.rstrip("/") self.system_prompt = system_prompt or _SYSTEM_PROMPT - self._conn: Optional[sqlite3.Connection] = None + self._conns = ThreadLocalConnections() # eine Connection pro Thread # Provider-spezifisch if provider == "anthropic": @@ -100,14 +103,10 @@ def __init__( # ollama braucht keine Extra-Init def _get_conn(self) -> sqlite3.Connection: - if self._conn is None: - self._conn = ensure_schema(self.knowledge_db) - return self._conn + return self._conns.get(self.knowledge_db) def close(self): - if self._conn: - self._conn.close() - self._conn = None + self._conns.close_all() # === Provider Backends === @@ -203,6 +202,9 @@ def summarize_queue(self, *, limit: int = 10, 'items': [], } + # Nach Absturz/Abbruch haengengebliebene Items wieder freigeben + stats['reset_stale'] = reset_stale_processing(conn) + queue_items = conn.execute(""" SELECT id, source_type, source_id FROM digest_queue @@ -248,10 +250,12 @@ def summarize_queue(self, *, limit: int = 10, 'output_tokens': 0, } + chunk_errors = [] for chunk_index, chunk_content in chunks: summary_result = self._summarize_chunk(chunk_content) if summary_result.get('error'): + chunk_errors.append((chunk_index, summary_result['error'])) continue conn.execute(""" @@ -285,6 +289,25 @@ def summarize_queue(self, *, limit: int = 10, if delay > 0: time.sleep(delay) + stats['total_input_tokens'] += item_result['input_tokens'] + stats['total_output_tokens'] += item_result['output_tokens'] + stats['items'].append(item_result) + + if chunk_errors: + # Mindestens ein Chunk fehlgeschlagen -> 'error' statt + # 'done' (erfolgreiche Summaries bleiben erhalten) + first_idx, first_err = chunk_errors[0] + item_result['chunk_errors'] = len(chunk_errors) + conn.execute( + "UPDATE digest_queue SET status='error', error_msg=?, " + "finished_at=CURRENT_TIMESTAMP WHERE id=?", + (f"{len(chunk_errors)}/{len(chunks)} Chunks fehlgeschlagen " + f"(Chunk {first_idx}: {first_err})"[:500], queue_id) + ) + conn.commit() + stats['errors'] += 1 + continue + conn.execute( "UPDATE digest_queue SET status='done', " "finished_at=CURRENT_TIMESTAMP WHERE id=?", @@ -293,9 +316,17 @@ def summarize_queue(self, *, limit: int = 10, conn.commit() stats['processed'] += 1 - stats['total_input_tokens'] += item_result['input_tokens'] - stats['total_output_tokens'] += item_result['output_tokens'] - stats['items'].append(item_result) + + except KeyboardInterrupt: + # Abbruch (Ctrl+C): aktuelles Item nicht in 'processing' + # haengen lassen, sondern wieder freigeben + conn.execute( + "UPDATE digest_queue SET status='pending', started_at=NULL " + "WHERE id=?", + (queue_id,) + ) + conn.commit() + raise except Exception as e: conn.execute( @@ -310,8 +341,9 @@ def summarize_queue(self, *, limit: int = 10, stats['duration_ms'] = elapsed if self.provider == "anthropic": - input_cost = stats['total_input_tokens'] / 1_000_000 * 0.25 - output_cost = stats['total_output_tokens'] / 1_000_000 * 1.25 + in_rate, out_rate = _ANTHROPIC_PRICE_PER_MTOK + input_cost = stats['total_input_tokens'] / 1_000_000 * in_rate + output_cost = stats['total_output_tokens'] / 1_000_000 * out_rate stats['estimated_cost_usd'] = round(input_cost + output_cost, 4) return stats @@ -466,8 +498,9 @@ def get_queue_status(self) -> Dict[str, Any]: } if self.provider == "anthropic": - input_cost = total_tokens_in / 1_000_000 * 0.25 - output_cost = total_tokens_out / 1_000_000 * 1.25 + in_rate, out_rate = _ANTHROPIC_PRICE_PER_MTOK + input_cost = total_tokens_in / 1_000_000 * in_rate + output_cost = total_tokens_out / 1_000_000 * out_rate result['summaries']['estimated_cost_usd'] = round( input_cost + output_cost, 4 ) diff --git a/tests/test_bugfixes_2026_10.py b/tests/test_bugfixes_2026_10.py new file mode 100644 index 0000000..b39561f --- /dev/null +++ b/tests/test_bugfixes_2026_10.py @@ -0,0 +1,738 @@ +""" +Regressionstests fuer die Bugfix-Runde 2026-10. + +Jeder Test fixiert einen verifizierten Befund (Nummern wie im Review): + 1 FTS5-Delete-Trigger ('delete'-Befehl auf regulaerer FTS5-Tabelle) + 2 Loeschen: erst DB (Transaktion), dann Datei verschieben + 3 Falsche relative Imports in gui/panels + 4 SQLite-Connection aus GUI-Thread in Scan-Threads + 5 EventBus ruft Handler im Worker-Thread auf + 6 Queue-Item 'done' obwohl Chunks fehlschlugen + 7 Queue-Items haengen nach Absturz in 'processing' + 8 Sortierte Dokumentliste -> falsches Dokument + 9 FTS5-Syntaxfehler durch Freitext / LIKE-Platzhalter +10 Re-Ingest: alte Summaries + Queue-Status +11 Archivierte Dokumente (archived_path) +12 Harte Chunk-Obergrenze +13 UTF-8-BOM verdeckt Frontmatter +14 Verzeichnisfilter matcht Geschwister-Ordner +15 is_active=0 wurde zu 1 + + Kostenschaetzung Haiku 4.5, zoll_station-Aufruf, open_path + +Keine Netzwerkzugriffe, keine echten LLM-Aufrufe. +""" + +import importlib +import json +import os +import sqlite3 +import subprocess +import sys +import threading +import time +import urllib.error +import urllib.request +from http.server import HTTPServer +from pathlib import Path + +import pytest + +_PKG_PARENT = str(Path(__file__).resolve().parent.parent.parent) +_PKG_NAME = Path(__file__).resolve().parent.parent.name + + +def _mod(name): + """Importiert ein Paket-Modul (digest.py & Co. nutzen relative Imports).""" + if _PKG_PARENT not in sys.path: + sys.path.insert(0, _PKG_PARENT) + pkg = _PKG_NAME if _PKG_NAME.isidentifier() else "KnowledgeDigest" + return importlib.import_module(f"{pkg}.{name}") + + +def _make_kd(tmp_path): + config_mod = _mod("config") + digest_mod = _mod("digest") + cfg = config_mod.Config(tmp_path / "cfg.json") + cfg.set("inbox_dir", str(tmp_path / "inbox")) + cfg.set("archive_dir", str(tmp_path / "archive")) + cfg.set("indexed_directories", []) + return digest_mod.KnowledgeDigest(db_path=tmp_path / "knowledge.db", config=cfg) + + +def _write(path, text): + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(text, encoding="utf-8") + return path + + +def _db(kd): + return _mod("schema").ensure_schema(kd.db_path) + + +def _qapp(): + pytest.importorskip("PySide6") + os.environ.setdefault("QT_QPA_PLATFORM", "offscreen") + from PySide6.QtWidgets import QApplication + return QApplication.instance() or QApplication([]) + + +# ============================================================ +# 1 FTS5-Delete-Trigger +# ============================================================ + +_LEGACY_DOC_TRIGGER = """ +CREATE TRIGGER doc_chunk_ad AFTER DELETE ON document_chunks BEGIN + INSERT INTO document_fts(document_fts, rowid, filename, content) + SELECT 'delete', old.id, + (SELECT filename FROM documents WHERE id = old.doc_id), + old.content; +END; +""" + + +class TestFtsDeleteTriggers: + def test_reingest_modified_file(self, tmp_path): + kd = _make_kd(tmp_path) + f = _write(tmp_path / "docs" / "a.txt", "Altinhalt Apfel. " * 100) + assert kd.ingest(f, archive=False)["status"] == "ok" + _write(f, "Neuinhalt Birne. " * 100) + result = kd.ingest(f, archive=False) + assert result["status"] == "ok", result["error"] + assert kd._search_documents("Birne") + assert not kd._search_documents("Apfel") + conn = _db(kd) + n_fts = conn.execute("SELECT COUNT(*) FROM document_fts").fetchone()[0] + n_chunks = conn.execute("SELECT COUNT(*) FROM document_chunks").fetchone()[0] + conn.close() + kd.close() + assert n_fts == n_chunks + + def test_legacy_trigger_is_migrated(self, tmp_path): + schema = _mod("schema") + db = tmp_path / "knowledge.db" + conn = schema.ensure_schema(db) + conn.execute("DROP TRIGGER doc_chunk_ad") + conn.executescript(_LEGACY_DOC_TRIGGER) + conn.execute("INSERT INTO documents (id, file_path, filename, file_type) " + "VALUES (1, 'x', 'x.txt', 'txt')") + conn.execute("INSERT INTO document_chunks (doc_id, chunk_index, content) " + "VALUES (1, 0, 'hallo welt')") + conn.commit() + # Vorbedingung: mit altem Trigger schlaegt jedes DELETE fehl + with pytest.raises(sqlite3.OperationalError): + conn.execute("DELETE FROM document_chunks WHERE doc_id = 1") + conn.rollback() + conn.close() + + conn = schema.ensure_schema(db) # Migration beim Oeffnen + sql = conn.execute("SELECT sql FROM sqlite_master WHERE name='doc_chunk_ad'").fetchone()[0] + assert "'delete'" not in sql + conn.execute("DELETE FROM document_chunks WHERE doc_id = 1") + conn.commit() + assert conn.execute("SELECT COUNT(*) FROM document_fts").fetchone()[0] == 0 + assert schema.get_schema_version(db) == schema.SCHEMA_VERSION == 5 + conn.close() + + def test_skill_and_wiki_chunk_delete(self, tmp_path): + schema = _mod("schema") + conn = schema.ensure_schema(tmp_path / "knowledge.db") + conn.execute("INSERT INTO skill_index (id, skill_name, skill_type) VALUES (1, 's', 't')") + conn.execute("INSERT INTO skill_chunks (skill_id, chunk_index, content) VALUES (1, 0, 'abc')") + conn.execute("INSERT INTO wiki_index (id, wiki_path, title) VALUES (1, 'w', 'W')") + conn.execute("INSERT INTO wiki_chunks (wiki_id, chunk_index, content) VALUES (1, 0, 'abc')") + conn.execute("DELETE FROM skill_chunks") + conn.execute("DELETE FROM wiki_chunks") + conn.commit() + assert conn.execute("SELECT COUNT(*) FROM skill_fts").fetchone()[0] == 0 + assert conn.execute("SELECT COUNT(*) FROM wiki_fts").fetchone()[0] == 0 + conn.close() + + +# ============================================================ +# 2 + 11 Loeschen (Web-Viewer + CLI deduplicate) +# ============================================================ + +class _Server: + def __init__(self, db_path): + wv = _mod("web_viewer") + handler = type("H", (wv.ViewerHandler,), {"db_path": Path(db_path)}) + self.httpd = HTTPServer(("127.0.0.1", 0), handler) + self.port = self.httpd.server_port + threading.Thread(target=self.httpd.serve_forever, daemon=True).start() + + def post(self, path): + req = urllib.request.Request(f"http://127.0.0.1:{self.port}{path}", + method="POST", data=b"") + try: + with urllib.request.urlopen(req, timeout=10) as resp: + return resp.status, json.loads(resp.read()) + except urllib.error.HTTPError as e: + return e.code, json.loads(e.read()) + + def stop(self): + self.httpd.shutdown() + self.httpd.server_close() + + +def _ingested_doc(tmp_path, name="a.txt"): + kd = _make_kd(tmp_path) + f = _write(tmp_path / "docs" / name, f"Inhalt {name}. " * 80) + assert kd.ingest(f, archive=False)["status"] == "ok" + kd.close() + return kd, f + + +class TestDeleteOrder: + def test_db_failure_leaves_file_in_place(self, tmp_path): + kd, f = _ingested_doc(tmp_path) + conn = _db(kd) + conn.execute("CREATE TRIGGER block_del BEFORE DELETE ON documents " + "BEGIN SELECT RAISE(ABORT, 'boom'); END;") + conn.commit() + srv = _Server(kd.db_path) + try: + status, body = srv.post("/api/delete_file/1") + finally: + srv.stop() + assert status == 500 and not body["ok"] + assert f.exists() # Datei nicht verschoben + # Transaktion zurueckgerollt: Chunks noch da + assert conn.execute("SELECT COUNT(*) FROM document_chunks").fetchone()[0] > 0 + conn.close() + + def test_move_failure_keeps_db_consistent(self, tmp_path, monkeypatch): + kd, f = _ingested_doc(tmp_path) + wv = _mod("web_viewer") + + def _fail(*_a, **_k): + raise PermissionError("gesperrt") + monkeypatch.setattr(wv, "move_to_trash", _fail) + srv = _Server(kd.db_path) + try: + status, body = srv.post("/api/delete_file/1") + finally: + srv.stop() + assert status == 200 and body["ok"] + assert "gesperrt" in body["warning"] + conn = _db(kd) + assert conn.execute("SELECT COUNT(*) FROM documents").fetchone()[0] == 0 + conn.close() + assert f.exists() + + def test_delete_archived_document_moves_archived_file(self, tmp_path): + kd, f = _ingested_doc(tmp_path) + archived = tmp_path / "archive" / "20260101_a.txt" + archived.parent.mkdir(parents=True) + f.rename(archived) + conn = _db(kd) + conn.execute("UPDATE documents SET archived_path=? WHERE id=1", (str(archived),)) + conn.commit() + conn.close() + srv = _Server(kd.db_path) + try: + status, body = srv.post("/api/delete_file/1") + finally: + srv.stop() + assert status == 200 and body["ok"] + assert not archived.exists() + assert (tmp_path / "_Papierkorb" / "20260101_a.txt").exists() + + def test_deduplicate_continues_after_item_error(self, tmp_path, monkeypatch, capsys): + digest_mod = _mod("digest") + kd = _make_kd(tmp_path) + files = [_write(tmp_path / "docs" / f"d{i}.txt", "gleich") for i in (1, 2, 3)] + conn = _db(kd) + for i, fp in enumerate(files, start=1): + conn.execute("INSERT INTO documents (id, file_path, filename, file_type, content_hash) " + "VALUES (?, ?, ?, 'txt', 'h')", (i, str(fp), fp.name)) + conn.execute("CREATE TRIGGER block2 BEFORE DELETE ON documents WHEN old.id = 2 " + "BEGIN SELECT RAISE(ABORT, 'boom'); END;") + conn.commit() + monkeypatch.setattr(digest_mod, "KnowledgeDigest", lambda *a, **k: kd) + monkeypatch.setattr(sys, "argv", ["knowledgedigest", "deduplicate", "-y"]) + assert digest_mod.main() == 0 + ids = [r[0] for r in conn.execute("SELECT id FROM documents ORDER BY id")] + conn.close() + assert ids == [1, 2] # 3 geloescht, 2 blockiert + assert files[0].exists() and files[1].exists() + assert not files[2].exists() + assert (tmp_path / "_Papierkorb" / "d3.txt").exists() + assert "boom" in capsys.readouterr().out + + +# ============================================================ +# 3 + 8 + 11 + 14 GUI-Panels +# ============================================================ + +class TestGuiPanels: + def test_preview_panel_imports_resolve(self, tmp_path): + _qapp() + kd, _f = _ingested_doc(tmp_path) + preview_mod = _mod("gui.panels.preview_panel") + panel = preview_mod.PreviewPanel() + panel._current_doc_db_path = lambda: kd.db_path + panel._show_chunks({"id": 1}) + assert "=== Chunk 0 ===" in panel._chunks_widget.toPlainText() + # Config-Import (vorher: from ..config -> ModuleNotFoundError) + assert isinstance(preview_mod.PreviewPanel._current_doc_db_path(panel), Path) + + def test_preview_uses_archived_path(self, tmp_path): + _qapp() + preview_mod = _mod("gui.panels.preview_panel") + archived = _write(tmp_path / "archive" / "x.txt", "archiviert") + panel = preview_mod.PreviewPanel() + panel._current_doc_db_path = lambda: tmp_path / "knowledge.db" + panel.show_document({"filename": "x.txt", "file_path": str(tmp_path / "weg.txt"), + "archived_path": str(archived)}) + assert panel._current_path == str(archived) + assert panel._text_widget.toPlainText() == "archiviert" + + def test_document_list_sorted_selection(self, tmp_path): + _qapp() + kd = _make_kd(tmp_path) + for name, n in (("alpha", 50), ("beta", 10), ("gamma", 30)): + kd.ingest(_write(tmp_path / "docs" / f"{name}.txt", f"{name} Inhalt hier. " * n), + archive=False) + dl_mod = _mod("gui.panels.document_list") + bus_mod = _mod("gui.event_bus") + selected = [] + handler = selected.append + bus = bus_mod.get_event_bus() + bus.subscribe(bus_mod.EventType.DOCUMENT_SELECTED, handler) + try: + panel = dl_mod.DocumentListPanel(kd) + assert panel._table.rowCount() == 3 # Import ...schema funktioniert + panel._table.sortItems(2) # nach Woertern + visual = [panel._table.item(r, 0).text() for r in range(3)] + assert visual == ["beta.txt", "gamma.txt", "alpha.txt"] + panel._table.setCurrentCell(0, 0) + assert selected[-1]["filename"] == "beta.txt" + panel._table.setCurrentCell(2, 0) + assert selected[-1]["filename"] == "alpha.txt" + finally: + bus.unsubscribe(bus_mod.EventType.DOCUMENT_SELECTED, handler) + kd.close() + + def test_document_list_directory_filter_excludes_siblings(self, tmp_path): + _qapp() + kd = _make_kd(tmp_path) + kd.ingest(_write(tmp_path / "docs" / "a.txt", "Erstes Dokument. " * 30), archive=False) + kd.ingest(_write(tmp_path / "docs" / "sub" / "b.txt", "Zweites Dokument. " * 30), + archive=False) + kd.ingest(_write(tmp_path / "docs2" / "c.txt", "Drittes Dokument. " * 30), archive=False) + dl_mod = _mod("gui.panels.document_list") + panel = dl_mod.DocumentListPanel(kd) + panel.filter_by_directory(str(tmp_path / "docs")) + names = sorted(d["filename"] for d in panel._docs) + kd.close() + assert names == ["a.txt", "b.txt"] + + +# ============================================================ +# 4 Thread-lokale Connections +# ============================================================ + +class TestThreadConnections: + def test_scan_in_worker_thread_after_gui_thread_use(self, tmp_path): + kd = _make_kd(tmp_path) + kd.get_status() # GUI-Thread oeffnet Connections (Ingestor, Summarizer) + src = tmp_path / "docs" + for i in range(3): + _write(src / f"f{i}.txt", f"Dokument Nummer {i} mit Text. " * 30) + out = {} + t = threading.Thread(target=lambda: out.update( + kd.scan_directory(str(src), archive=False))) + t.start() + t.join() + assert out["errors"] == 0, [f["error"] for f in out["files"]] + assert out["ingested"] == 3 + assert kd.get_status()["documents"]["total_documents"] == 3 + kd.close() + + def test_close_all_and_reopen(self, tmp_path): + schema = _mod("schema") + pool = schema.ThreadLocalConnections() + db = tmp_path / "knowledge.db" + c1 = pool.get(db) + assert pool.get(db) is c1 + other = {} + t = threading.Thread(target=lambda: other.update(c=pool.get(db))) + t.start() + t.join() + assert other["c"] is not c1 + pool.close_all() + with pytest.raises(sqlite3.ProgrammingError): + c1.execute("SELECT 1") + with pytest.raises(sqlite3.ProgrammingError): + other["c"].execute("SELECT 1") + c2 = pool.get(db) + assert c2 is not c1 + c2.execute("SELECT 1") + pool.close_all() + + +# ============================================================ +# 5 EventBus: Handler im GUI-Thread +# ============================================================ + +class TestEventBusThreading: + def test_emit_from_worker_dispatches_in_gui_thread(self): + app = _qapp() + bus_mod = _mod("gui.event_bus") + bus = bus_mod.EventBus() + calls = [] + bus.subscribe(bus_mod.EventType.STATUS_MESSAGE, + lambda data: calls.append((data, threading.current_thread()))) + + t = threading.Thread(target=lambda: bus.emit(bus_mod.EventType.STATUS_MESSAGE, "aus Thread")) + t.start() + t.join() + assert calls == [] # nicht synchron im Worker-Thread ausgefuehrt + deadline = time.time() + 5 + while not calls and time.time() < deadline: + app.processEvents() + assert calls == [("aus Thread", threading.main_thread())] + + def test_emit_in_gui_thread_is_synchronous(self): + _qapp() + bus_mod = _mod("gui.event_bus") + bus = bus_mod.EventBus() + calls = [] + bus.subscribe(bus_mod.EventType.STATUS_MESSAGE, calls.append) + bus.emit(bus_mod.EventType.STATUS_MESSAGE, "direkt") + assert calls == ["direkt"] + + +# ============================================================ +# 6 + 7 Summarizer-Queue +# ============================================================ + +def _queue_db(tmp_path, n_chunks=2): + schema = _mod("schema") + db = tmp_path / "knowledge.db" + conn = schema.ensure_schema(db) + conn.execute("INSERT INTO documents (id, file_path, filename, file_type) " + "VALUES (1, 'x', 'x.txt', 'txt')") + for i in range(n_chunks): + conn.execute("INSERT INTO document_chunks (doc_id, chunk_index, content) " + "VALUES (1, ?, ?)", (i, f"chunk {i}")) + conn.execute("INSERT INTO digest_queue (source_type, source_id) VALUES ('document', 1)") + conn.commit() + return db, conn + + +_OK_JSON = '{"summary": "s", "keywords": ["k"], "domain": "d"}' + + +class TestSummarizerQueue: + def test_partial_chunk_failure_marks_error(self, tmp_path): + Summarizer = _mod("summarizer").Summarizer + db, conn = _queue_db(tmp_path) + + def llm(_system, text): + if text == "chunk 1": + raise RuntimeError("API down / 429") + return _OK_JSON + s = Summarizer(db, provider="custom", llm_fn=llm) + stats = s.summarize_queue(limit=5, delay=0) + s.close() + row = conn.execute("SELECT status, error_msg FROM digest_queue").fetchone() + assert (stats["processed"], stats["errors"]) == (0, 1) + assert row["status"] == "error" + assert "429" in row["error_msg"] and "1/2" in row["error_msg"] + # Erfolgreiche Summary bleibt erhalten + assert conn.execute("SELECT chunk_index FROM summaries").fetchall()[0][0] == 0 + conn.close() + + def test_all_chunks_ok_marks_done(self, tmp_path): + Summarizer = _mod("summarizer").Summarizer + db, conn = _queue_db(tmp_path) + s = Summarizer(db, provider="custom", llm_fn=lambda _s, _t: _OK_JSON) + stats = s.summarize_queue(limit=5, delay=0) + s.close() + assert (stats["processed"], stats["errors"]) == (1, 0) + assert conn.execute("SELECT status FROM digest_queue").fetchone()[0] == "done" + conn.close() + + def test_gemini_chunk_failure_marks_error(self, tmp_path): + gfs = _mod("gemini_flash_summarizer") + db, conn = _queue_db(tmp_path, n_chunks=1) + s = gfs.GeminiFlashSummarizer(db, api_key="k") + s._summarize_chunk = lambda _text: {"error": "quota"} + stats = s.summarize_queue(limit=5, delay=0) + s.close() + assert stats["errors"] == 1 and stats["processed"] == 0 + row = conn.execute("SELECT status, error_msg FROM digest_queue").fetchone() + assert row["status"] == "error" and "quota" in row["error_msg"] + conn.close() + + def test_stale_processing_is_reset(self, tmp_path): + Summarizer = _mod("summarizer").Summarizer + db, conn = _queue_db(tmp_path, n_chunks=1) + conn.execute("INSERT INTO documents (id, file_path, filename, file_type) " + "VALUES (2, 'y', 'y.txt', 'txt'), (3, 'z', 'z.txt', 'txt')") + conn.execute("INSERT INTO document_chunks (doc_id, chunk_index, content) " + "VALUES (2, 0, 'a'), (3, 0, 'b')") + conn.execute("UPDATE digest_queue SET status='processing', " + "started_at=datetime('now', '-2 hours') WHERE source_id=1") + conn.execute("INSERT INTO digest_queue (source_type, source_id, status, started_at) " + "VALUES ('document', 2, 'processing', NULL), " + "('document', 3, 'processing', CURRENT_TIMESTAMP)") + conn.commit() + s = Summarizer(db, provider="custom", llm_fn=lambda _s, _t: _OK_JSON) + stats = s.summarize_queue(limit=10, delay=0) + s.close() + status = dict(conn.execute("SELECT source_id, status FROM digest_queue").fetchall()) + conn.close() + assert stats["reset_stale"] == 2 + assert status == {1: "done", 2: "done", 3: "processing"} # 3 laeuft noch + + def test_keyboard_interrupt_resets_item_to_pending(self, tmp_path): + Summarizer = _mod("summarizer").Summarizer + db, conn = _queue_db(tmp_path) + + def llm(_system, _text): + raise KeyboardInterrupt + s = Summarizer(db, provider="custom", llm_fn=llm) + with pytest.raises(KeyboardInterrupt): + s.summarize_queue(limit=5, delay=0) + s.close() + row = conn.execute("SELECT status, started_at FROM digest_queue").fetchone() + conn.close() + assert row["status"] == "pending" and row["started_at"] is None + + def test_haiku_cost_estimate(self, tmp_path): + Summarizer = _mod("summarizer").Summarizer + db, conn = _queue_db(tmp_path, n_chunks=1) + conn.execute("INSERT INTO summaries (source_type, source_id, chunk_index, summary, " + "input_tokens, output_tokens) VALUES ('document', 1, 0, 's', 1000000, 1000000)") + conn.commit() + conn.close() + s = Summarizer(db, provider="anthropic", api_key="unused") + cost = s.get_queue_status()["summaries"]["estimated_cost_usd"] + s.close() + assert cost == pytest.approx(6.0) # $1 in + $5 out pro MTok + + +# ============================================================ +# 9 Suche: FTS5-Syntaxfehler + LIKE-Platzhalter +# ============================================================ + +class TestSearchSanitizing: + @pytest.fixture + def kd(self, tmp_path): + kd = _make_kd(tmp_path) + kd.ingest(_write(tmp_path / "docs" / "covid.txt", + "Studie zu COVID-19 und E-Mail Verkehr mit c++ Code. " * 20), + archive=False) + kd.ingest(_write(tmp_path / "docs" / "rabatt.txt", "Heute gibt es 50% Rabatt. " * 20), + archive=False) + yield kd + kd.close() + + @pytest.mark.parametrize("query", ["COVID-19", "E-Mail", "c++", '"COVID-19', "e-mail verkehr"]) + def test_free_text_queries_find_document(self, kd, query): + names = [r["name"] for r in kd.search_all(query)] + assert "covid.txt" in names + + def test_fts_syntax_still_works(self, kd): + assert [r["name"] for r in kd.search_all("Studie AND Verkehr")] == ["covid.txt"] + + def test_document_like_fallback_escapes_wildcards(self, kd): + conn = kd._get_conn() + try: + assert kd._search_documents_fallback(conn, "%%", 10) == [] + assert kd._search_documents_fallback(conn, "5_%", 10) == [] + hits = kd._search_documents_fallback(conn, "50%", 10) + finally: + conn.close() + assert [h["name"] for h in hits] == ["rabatt.txt"] + + def test_skill_like_fallback_escapes_wildcards(self, kd): + conn = kd._get_conn() + conn.execute("INSERT INTO skill_index (id, skill_name, skill_type) VALUES (1, 'abc', 't')") + conn.commit() + try: + assert kd._search_fallback(conn, "%", 10, None, None) == [] + assert kd._search_wikis_fallback(conn, "_", 10, None) == [] + finally: + conn.close() + + def test_web_viewer_search(self, kd): + wv = _mod("web_viewer") + assert "covid.txt" in wv.page_search(kd.db_path, query="COVID-19") + # '%' ist kein LIKE-Platzhalter mehr (vorher: Treffer auf alles) + assert "covid.txt" not in wv.page_search(kd.db_path, query="%") + assert "covid.txt" not in wv.page_search(kd.db_path, query='"_') + + +# ============================================================ +# 10 Re-Ingest setzt Summaries + Queue zurueck +# ============================================================ + +def test_reingest_resets_summaries_and_queue(tmp_path): + kd = _make_kd(tmp_path) + f = _write(tmp_path / "docs" / "a.txt", "Version eins. " * 80) + kd.ingest(f, archive=False) + conn = _db(kd) + conn.execute("INSERT INTO summaries (source_type, source_id, chunk_index, summary) " + "VALUES ('document', 1, 0, 'alt')") + conn.execute("UPDATE digest_queue SET status='done', finished_at=CURRENT_TIMESTAMP, " + "error_msg='x'") + conn.commit() + _write(f, "Version zwei ganz anders. " * 80) + assert kd.ingest(f, archive=False)["status"] == "ok" + kd.close() + assert conn.execute("SELECT COUNT(*) FROM summaries").fetchone()[0] == 0 + row = conn.execute("SELECT status, finished_at, error_msg FROM digest_queue").fetchone() + conn.close() + assert tuple(row) == ("pending", None, None) + + +# ============================================================ +# 11 resolve_document_path +# ============================================================ + +def test_resolve_document_path(tmp_path): + utils = _mod("utils") + present = _write(tmp_path / "da.txt", "x") + archived = _write(tmp_path / "archiv.txt", "x") + missing = str(tmp_path / "weg.txt") + assert utils.resolve_document_path({"file_path": str(present), + "archived_path": str(archived)}) == str(present) + assert utils.resolve_document_path({"file_path": missing, + "archived_path": str(archived)}) == str(archived) + assert utils.resolve_document_path({"file_path": missing}) == missing + # sqlite3.Row ohne archived_path-Spalte + conn = sqlite3.connect(":memory:") + conn.row_factory = sqlite3.Row + row = conn.execute("SELECT ? AS file_path", (missing,)).fetchone() + assert utils.resolve_document_path(row) == missing + + +# ============================================================ +# 12 + 13 Chunker / Extractor +# ============================================================ + +def test_chunker_enforces_hard_upper_bound(): + chunker = _mod("chunker") + text = " ".join(f"wort{i}" for i in range(3000)) + chunks = chunker.chunk_text(text) + assert all(c.token_count <= chunker.MAX_CHUNK_SIZE for c in chunks) + assert " ".join(c.content for c in chunks).split() == text.split() + assert [c.index for c in chunks] == list(range(len(chunks))) + + +def test_bom_text_file_frontmatter_detected(tmp_path): + extractor = _mod("extractor") + chunker = _mod("chunker") + f = tmp_path / "skill.md" + f.write_bytes("---\ntitle: x\n---\n# Body\n".encode("utf-8-sig") + b"text " * 100) + text = extractor.TextExtractor().extract(f).text + assert not text.startswith("\ufeff") + assert chunker.chunk_text(text)[0].is_frontmatter + # auch direkt uebergebener Text mit BOM + assert chunker.chunk_text("\ufeff---\ntitle: x\n---\n# Body\n" + "text " * 100)[0].is_frontmatter + + +# ============================================================ +# 14 Verzeichnisfilter +# ============================================================ + +class TestDirectoryFilter: + def _insert(self, conn, source_dirs): + for i, d in enumerate(source_dirs, start=1): + conn.execute("INSERT INTO documents (file_path, filename, file_type, source_dir) " + "VALUES (?, ?, 'txt', ?)", (f"{d}/f{i}", f"f{i}", d)) + conn.commit() + + def test_get_directories_counts(self, tmp_path): + kd = _make_kd(tmp_path) + base = str(tmp_path / "a_b") + conn = _db(kd) + self._insert(conn, [base, base + "/sub", base + "\\win", base + "2", + str(tmp_path / "aXb"), base + "%x"]) + conn.close() + kd._config.set("indexed_directories", [base]) + assert kd.get_directories() == [{"path": base, "doc_count": 3}] + + def test_dir_filter_sql_trailing_separator(self, tmp_path): + utils = _mod("utils") + conn = sqlite3.connect(":memory:") + conn.execute("CREATE TABLE documents (source_dir TEXT)") + conn.executemany("INSERT INTO documents VALUES (?)", + [("/data/docs",), ("/data/docs/x",), ("/data/docs2",)]) + sql, params = utils.dir_filter_sql("source_dir", "/data/docs/") + n = conn.execute(f"SELECT COUNT(*) FROM documents WHERE {sql}", params).fetchone()[0] + assert n == 2 + + +# ============================================================ +# 15 is_active=0 bleibt 0 +# ============================================================ + +def test_skill_indexer_keeps_inactive(tmp_path): + indexer = _mod("indexer") + bach = tmp_path / "bach.db" + conn = sqlite3.connect(str(bach)) + conn.execute("CREATE TABLE skills (id INTEGER, name TEXT, type TEXT, category TEXT, " + "path TEXT, description TEXT, version TEXT, content TEXT, " + "content_hash TEXT, is_active INTEGER)") + conn.executemany( + "INSERT INTO skills VALUES (?, ?, 'agent', 'c', 'p', 'd', '1', ?, NULL, ?)", + [(1, "aktiv", "Inhalt eins " * 20, 1), + (2, "inaktiv", "Inhalt zwei " * 20, 0), + (3, "unbekannt", "Inhalt drei " * 20, None)]) + conn.commit() + conn.close() + idx = indexer.SkillIndexer(tmp_path / "knowledge.db") + idx.index_from_bach(bach) + rows = dict(idx._get_conn().execute("SELECT skill_name, is_active FROM skill_index").fetchall()) + idx.close() + assert rows == {"aktiv": 1, "inaktiv": 0, "unbekannt": 1} + + +# ============================================================ +# Kleinere Befunde: zoll_station, open_path +# ============================================================ + +def test_zoll_station_runs_digest_as_module(): + zoll = _mod("zoll_station") + pkg_dir = Path(zoll.__file__).resolve().parent + cmd, cwd = zoll._digest_command(pkg_dir, "summarize", "--help") + assert cmd[:2] == [sys.executable, "-m"] + assert "digest.py" not in " ".join(cmd) + assert cwd == str(pkg_dir.parent) + result = subprocess.run(cmd, cwd=cwd, capture_output=True, text=True, timeout=60, check=False, + env={**os.environ, "PYTHONIOENCODING": "utf-8"}) + assert result.returncode == 0, result.stderr + assert "--flash" in result.stdout + + +@pytest.mark.parametrize("platform, expected", [("darwin", "open"), ("linux", "xdg-open")]) +def test_open_path_cross_platform(monkeypatch, platform, expected): + utils = _mod("utils") + calls = [] + monkeypatch.setattr(utils.sys, "platform", platform) + monkeypatch.setattr(utils.subprocess, "Popen", lambda args: calls.append(args)) + utils.open_path(Path("/tmp/x.pdf")) + assert calls == [[expected, str(Path("/tmp/x.pdf"))]] + + +def test_open_path_windows(monkeypatch): + utils = _mod("utils") + calls = [] + monkeypatch.setattr(utils.sys, "platform", "win32") + monkeypatch.setattr(utils.os, "startfile", calls.append, raising=False) + utils.open_path("C:\\x.pdf") + assert calls == ["C:\\x.pdf"] + + +def test_config_defaults_are_not_shared_between_instances(tmp_path): + from KnowledgeDigest.config import DEFAULT_CONFIG, Config + + before = list(DEFAULT_CONFIG["indexed_directories"]) + cfg = Config(tmp_path / "config.json") + cfg._data["indexed_directories"].append(str(tmp_path)) + assert DEFAULT_CONFIG["indexed_directories"] == before + assert Config(tmp_path / "other.json")._data["indexed_directories"] == before diff --git a/tests/test_core.py b/tests/test_core.py index 7513a15..4d5f4ec 100644 --- a/tests/test_core.py +++ b/tests/test_core.py @@ -188,12 +188,12 @@ def test_large_text_multiple_chunks(self): assert len(chunks) >= 3 def test_single_line_no_paragraphs_exceeds_max(self): - # Ein einziger Einzeiler > MAX_CHUNK_SIZE (500) wird als ein oversized Chunk - # unveraendert zurueckgegeben (kein Split ohne Satzzeichen/Absaetze moeglich) + # Ein einziger Einzeiler > MAX_CHUNK_SIZE (500) ohne Satzzeichen/Absaetze + # wird hart in Wortfenster der Groesse chunk_size zerlegt (harte Obergrenze) text = " ".join([f"word{i}" for i in range(1400)]) chunks = chunk_text(text, chunk_size=350) - assert len(chunks) == 1 - assert chunks[0].token_count == 1400 + assert [c.token_count for c in chunks] == [350, 350, 350, 350] + assert " ".join(c.content for c in chunks) == text def test_token_count_matches_content(self): text = "Dies ist ein Test mit genau acht Woertern hier." @@ -276,7 +276,7 @@ def test_schema_version_written(self): db_path.unlink(missing_ok=True) def test_schema_version_constant(self): - assert SCHEMA_VERSION == 4 + assert SCHEMA_VERSION == 5 def test_get_schema_version_nonexistent(self, tmp_path): # Datei existiert nicht -> Version 0 diff --git a/tests/test_transit.py b/tests/test_transit.py index a11aca7..26f31cf 100644 --- a/tests/test_transit.py +++ b/tests/test_transit.py @@ -15,10 +15,15 @@ from datetime import timedelta import pytest -from sqlite_transit_sync import SyncError + +try: + from sqlite_transit_sync import SyncError +except ImportError: # optional dependency, see transit.py + SyncError = None from KnowledgeDigest.schema import ensure_schema from KnowledgeDigest.transit import ( + TRANSIT_SYNC_AVAILABLE, generate_chunk_key, ChunkTaskManager, KnowledgeDigestMergePolicy, @@ -27,6 +32,10 @@ iso_timestamp, ) +requires_transit_sync = pytest.mark.skipif( + not TRANSIT_SYNC_AVAILABLE, reason="optional package sqlite-transit-sync not installed" +) + @pytest.fixture def transit_environment(tmp_path): @@ -86,6 +95,7 @@ def test_chunk_key_generation(): assert key1 != key_modified +@requires_transit_sync def test_t01_two_nodes_exclusive_claim(transit_environment): """Criterion 1: Both nodes see the same chunk, but only one node obtains the claim.""" env = transit_environment @@ -143,6 +153,7 @@ def test_t01_two_nodes_exclusive_claim(transit_environment): assert started is True +@requires_transit_sync def test_t02_result_converges_once(transit_environment): """Criterion 2: Result appears after pull on both nodes exactly once.""" env = transit_environment @@ -188,6 +199,7 @@ def test_t02_result_converges_once(transit_environment): assert claim_attempt is False +@requires_transit_sync def test_t03_expired_claim_takeover(transit_environment): """Criterion 3: Expired claim can be safely taken over by another node.""" env = transit_environment @@ -272,6 +284,7 @@ def test_t03_expired_claim_takeover(transit_environment): assert task_a_final["claimed_by_node"] == "node_b" +@requires_transit_sync def test_t04_repeated_pull_is_idempotent(transit_environment): """Criterion 4: Repeated pull operations are strictly idempotent.""" env = transit_environment @@ -314,6 +327,7 @@ def test_t04_repeated_pull_is_idempotent(transit_environment): assert report.unchanged == 5 +@requires_transit_sync def test_t05_secret_negative_test_and_quick_check(transit_environment): """Criterion 5: Secret negative test blocks export, and PRAGMA quick_check stays ok.""" env = transit_environment @@ -373,3 +387,105 @@ def test_t05_secret_negative_test_and_quick_check(transit_environment): # Quick check remains ok check_a_after = conn_a.execute("PRAGMA quick_check").fetchone()[0] assert check_a_after == "ok" + + +# --------------------------------------------------------------------------- +# Lifecycle tests that run without the optional sqlite-transit-sync package +# --------------------------------------------------------------------------- + + +@pytest.fixture +def single_node(tmp_path): + conn = ensure_schema(tmp_path / "knowledge.db") + yield conn + conn.close() + + +def test_iso_timestamp_is_fixed_width_utc(): + from datetime import datetime, timezone, timedelta as td + + whole = iso_timestamp(datetime(2026, 1, 1, 12, 0, 0, tzinfo=timezone.utc)) + frac = iso_timestamp(datetime(2026, 1, 1, 12, 0, 0, 500, tzinfo=timezone.utc)) + shifted = iso_timestamp(datetime(2026, 1, 1, 14, 0, 0, tzinfo=timezone(td(hours=2)))) + assert len(whole) == len(frac) + assert whole < frac + assert shifted == whole + + +def test_register_task_is_idempotent(single_node): + key1 = ChunkTaskManager.register_task(single_node, "document", "a.md", "v1", 0) + key2 = ChunkTaskManager.register_task(single_node, "document", "a.md", "v1", 0) + assert key1 == key2 + count = single_node.execute("SELECT COUNT(*) FROM chunk_tasks").fetchone()[0] + assert count == 1 + + +def test_lifecycle_claim_process_heartbeat_done(single_node): + key = ChunkTaskManager.register_task(single_node, "document", "b.md", "v1", 0) + assert ChunkTaskManager.claim_task(single_node, key, "n1", "agent") is True + # second claim on an active lease is rejected + assert ChunkTaskManager.claim_task(single_node, key, "n2", "agent") is False + # only the owning node may start processing + assert ChunkTaskManager.start_processing(single_node, key, "n2") is False + assert ChunkTaskManager.start_processing(single_node, key, "n1") is True + assert ChunkTaskManager.heartbeat(single_node, key, "n1") is True + assert ChunkTaskManager.complete_task(single_node, key, "n1", "summary") is True + task = ChunkTaskManager.get_task(single_node, key) + assert task["status"] == "done" + # done is terminal + assert ChunkTaskManager.claim_task(single_node, key, "n2", "agent") is False + assert ChunkTaskManager.heartbeat(single_node, key, "n1") is False + + +def test_get_task_does_not_change_row_factory(single_node): + before = single_node.row_factory + key = ChunkTaskManager.register_task(single_node, "wiki", "Page", "v1", 0) + assert ChunkTaskManager.get_task(single_node, key)["chunk_key"] == key + assert single_node.row_factory is before + + +def test_failed_task_can_be_reclaimed(single_node): + key = ChunkTaskManager.register_task(single_node, "skill", "s", "v1", 0) + ChunkTaskManager.claim_task(single_node, key, "n1", "agent") + assert ChunkTaskManager.fail_task(single_node, key, "n1", "boom") is True + assert ChunkTaskManager.get_task(single_node, key)["status"] == "error" + assert ChunkTaskManager.claim_task(single_node, key, "n2", "agent") is True + + +def test_expired_lease_takeover_without_sync(single_node): + key = ChunkTaskManager.register_task(single_node, "document", "c.md", "v1", 0) + past = utc_now() - timedelta(minutes=10) + assert ChunkTaskManager.claim_task(single_node, key, "n1", "agent", lease_seconds=60, now=past) + assert ChunkTaskManager.claim_task(single_node, key, "n2", "agent") is True + assert ChunkTaskManager.get_task(single_node, key)["claimed_by_node"] == "n2" + + +@pytest.mark.parametrize( + "local, remote, expected", + [ + ({"status": "done"}, {"status": "claimed"}, False), + ({"status": "claimed", "lease_expires_at": "9999"}, {"status": "done"}, True), + ({"status": "pending"}, {"status": "error"}, True), + ({"status": "pending"}, {"status": "claimed", "lease_expires_at": "9999"}, True), + ( + {"status": "claimed", "lease_expires_at": "9999", "claimed_at": "2026-01-02"}, + {"status": "claimed", "lease_expires_at": "9999", "claimed_at": "2026-01-01"}, + True, + ), + ( + {"status": "claimed", "lease_expires_at": "9999", "claimed_at": "x", "claimed_by_node": "a"}, + {"status": "claimed", "lease_expires_at": "9999", "claimed_at": "x", "claimed_by_node": "b"}, + False, + ), + ], +) +def test_resolve_chunk_conflict(local, remote, expected): + remote_wins, _ = KnowledgeDigestMergePolicy._resolve_chunk_conflict(local, remote, iso_timestamp()) + assert remote_wins is expected + + +def test_create_knowledge_sync_reports_missing_dependency(tmp_path): + if TRANSIT_SYNC_AVAILABLE: + pytest.skip("sqlite-transit-sync installed") + with pytest.raises(ImportError, match="sqlite-transit-sync"): + create_knowledge_sync(tmp_path / "k.db", tmp_path / "yard", "n1") diff --git a/transit.py b/transit.py index 9c0e05c..04ab8aa 100644 --- a/transit.py +++ b/transit.py @@ -28,23 +28,46 @@ from pathlib import Path from typing import Any, Optional -from sqlite_transit_sync import ( - SyncConfig, - TransitSync, - Snapshot, - MergeReport, +try: # optional dependency: https://github.com/ellmos-ai/sqlite-transit-sync + from sqlite_transit_sync import ( + SyncConfig, + TransitSync, + Snapshot, + MergeReport, + ) + + TRANSIT_SYNC_AVAILABLE = True +except ImportError: # pragma: no cover - depends on the local environment + SyncConfig = TransitSync = Snapshot = MergeReport = None # type: ignore[assignment,misc] + TRANSIT_SYNC_AVAILABLE = False + +_MISSING_TRANSIT_MSG = ( + "Transit sync requires the optional package 'sqlite-transit-sync' " + "(https://github.com/ellmos-ai/sqlite-transit-sync). " + "ChunkTaskManager and generate_chunk_key work without it." ) +def _require_transit_sync() -> None: + if not TRANSIT_SYNC_AVAILABLE: + raise ImportError(_MISSING_TRANSIT_MSG) + + def utc_now() -> datetime: """Returns current UTC timestamp with timezone.""" return datetime.now(timezone.utc) def iso_timestamp(dt: Optional[datetime] = None) -> str: - """Formats a datetime as ISO 8601 string.""" + """Formats a datetime as fixed-width ISO 8601 UTC string. + + Lease checks compare timestamps lexicographically (in SQL and in the merge + policy), so every value must use the same offset and precision. + """ t = dt or utc_now() - return t.isoformat() + if t.tzinfo is None: + t = t.replace(tzinfo=timezone.utc) + return t.astimezone(timezone.utc).isoformat(timespec="microseconds") def parse_iso(ts_str: Optional[str]) -> Optional[datetime]: @@ -278,9 +301,9 @@ def fail_task( @staticmethod def get_task(conn: sqlite3.Connection, chunk_key: str) -> Optional[dict[str, Any]]: """Fetches a task by chunk_key as dict.""" - conn.row_factory = sqlite3.Row - cur = conn.execute("SELECT * FROM chunk_tasks WHERE chunk_key = ?", (chunk_key,)) - row = cur.fetchone() + cur = conn.cursor() + cur.row_factory = sqlite3.Row # do not mutate the caller's connection + row = cur.execute("SELECT * FROM chunk_tasks WHERE chunk_key = ?", (chunk_key,)).fetchone() return dict(row) if row else None @@ -305,6 +328,7 @@ def merge( remote: sqlite3.Connection, snapshot: Snapshot, ) -> MergeReport: + _require_transit_sync() local.row_factory = sqlite3.Row remote.row_factory = sqlite3.Row @@ -497,6 +521,7 @@ def create_knowledge_sync( state_dir: Directory for local sync state (defaults to ~/.wissensdb-transit). namespace: Logical namespace for transit snapshots. """ + _require_transit_sync() if state_dir is None: state_dir = db_path.parent / ".transit_state" state_dir.mkdir(parents=True, exist_ok=True) diff --git a/utils.py b/utils.py index e300af0..4feee1e 100644 --- a/utils.py +++ b/utils.py @@ -6,13 +6,26 @@ - sha256_hash(): Content-Hashing fuer Aenderungserkennung - extract_keywords(): Keyword-Extraktion mit Stoppwort-Filterung - STOP_WORDS: DE/EN Stoppwoerter + - resolve_document_path(): aktueller Pfad eines Dokuments (auch archiviert) + - open_path(): Datei/Ordner mit dem Standardprogramm oeffnen (plattformuebergreifend) + - escape_like(): Platzhalter % und _ fuer LIKE ... ESCAPE '\\' maskieren + - dir_filter_sql(): SQL-Filter "liegt in Verzeichnis X oder darunter" + - fts5_quote(): Freitext-Suche als sichere FTS5-Query (Tokens in Anfuehrungszeichen) + - move_to_trash(): Datei kollisionsfrei in einen Papierkorb-Ordner verschieben """ -__all__ = ["sha256_hash", "extract_keywords", "STOP_WORDS"] +__all__ = ["sha256_hash", "extract_keywords", "STOP_WORDS", + "resolve_document_path", "open_path", "escape_like", + "dir_filter_sql", "fts5_quote", "move_to_trash"] import hashlib +import os import re -from typing import List, Dict +import shutil +import subprocess +import sys +from pathlib import Path +from typing import Any, List, Dict, Optional, Tuple # Stoppwoerter (DE/EN gemischt, fuer Keyword-Extraktion) @@ -76,3 +89,93 @@ def extract_keywords(text: str, max_count: int = MAX_KEYWORDS_PER_ITEM) -> List[ # Nach Haeufigkeit sortieren, Top-N sorted_kw = sorted(freq.items(), key=lambda x: x[1], reverse=True) return [kw for kw, _ in sorted_kw[:max_count]] + + +def resolve_document_path(doc: Any) -> Optional[str]: + """Gibt den aktuell gueltigen Pfad eines Dokuments zurueck. + + Archivierte Dokumente liegen nicht mehr unter file_path, sondern unter + archived_path. Reihenfolge: file_path (falls vorhanden), sonst + archived_path (falls vorhanden), sonst file_path als Fallback. + + Args: + doc: dict oder sqlite3.Row mit 'file_path' und optional 'archived_path' + """ + def _get(key: str) -> Optional[str]: + try: + return doc[key] + except (KeyError, IndexError): + return None + + file_path = _get("file_path") + archived_path = _get("archived_path") + if file_path and os.path.exists(file_path): + return file_path + if archived_path and os.path.exists(archived_path): + return archived_path + return file_path or archived_path or None + + +def open_path(path: Any) -> None: + """Oeffnet Datei/Ordner mit dem Standardprogramm des Systems. + + Windows: os.startfile, macOS: open, Linux/sonstige: xdg-open. + """ + path = str(path) + if sys.platform == "win32": + os.startfile(path) # type: ignore[attr-defined] + elif sys.platform == "darwin": + subprocess.Popen(["open", path]) + else: + subprocess.Popen(["xdg-open", path]) + + +def escape_like(text: str) -> str: + """Maskiert \\, % und _ fuer SQL ``LIKE ? ESCAPE '\\'``.""" + return (text.replace("\\", "\\\\") + .replace("%", "\\%") + .replace("_", "\\_")) + + +def dir_filter_sql(column: str, directory: str) -> Tuple[str, List[str]]: + """SQL-Bedingung: `column` ist `directory` selbst oder liegt darunter. + + Anders als ``LIKE directory || '%'`` matcht das keine Geschwister-Ordner + (``docs`` vs. ``docs2``) und behandelt % / _ im Pfad nicht als Platzhalter. + Beide Trennzeichen (/ und \\) werden beruecksichtigt. + + Returns: + (sql, params) -- z.B. ``("(d.source_dir = ? OR ...)", [...])`` + """ + base = directory.rstrip("/\\") + esc = escape_like(base) + sql = (f"({column} = ? OR {column} LIKE ? ESCAPE '\\' " + f"OR {column} LIKE ? ESCAPE '\\')") + return sql, [base, esc + "/%", esc + "\\\\%"] + + +def fts5_quote(query: str) -> str: + """Macht aus Freitext eine syntaktisch sichere FTS5-Query. + + Jedes Whitespace-Token wird in Anfuehrungszeichen gesetzt (innere " werden + verdoppelt), z.B. ``COVID-19 c++`` -> ``"COVID-19" "c++"``. + """ + return " ".join('"' + tok.replace('"', '""') + '"' for tok in query.split()) + + +def move_to_trash(file_path: Any, trash_dir: Any) -> Path: + """Verschiebt eine Datei in trash_dir (Name bei Kollision mit _1, _2, ...). + + Gibt den Zielpfad zurueck; Fehler (z.B. Datei gesperrt) werden geworfen. + """ + trash_dir = Path(trash_dir) + trash_dir.mkdir(parents=True, exist_ok=True) + base_name = os.path.basename(str(file_path)) + name, ext = os.path.splitext(base_name) + trash_path = trash_dir / base_name + counter = 1 + while trash_path.exists(): + trash_path = trash_dir / f"{name}_{counter}{ext}" + counter += 1 + shutil.move(str(file_path), str(trash_path)) + return trash_path diff --git a/web_viewer.py b/web_viewer.py index a208640..3122558 100644 --- a/web_viewer.py +++ b/web_viewer.py @@ -15,7 +15,6 @@ import sys import json import html -import shutil import sqlite3 import argparse import webbrowser @@ -24,6 +23,11 @@ from http.server import HTTPServer, BaseHTTPRequestHandler from pathlib import Path +try: + from .utils import escape_like, fts5_quote, move_to_trash, open_path, resolve_document_path +except ImportError: # Standalone (python web_viewer.py / Tests mit sys.path) + from utils import escape_like, fts5_quote, move_to_trash, open_path, resolve_document_path + def _get_conn(db_path): conn = sqlite3.connect(str(db_path), timeout=30) @@ -323,7 +327,10 @@ def page_doc(db_path, doc_id): fetch('/api/delete_file/' + docId, {{method: 'POST'}}) .then(r => r.json()) .then(data => {{ - if(data.ok) location.href = '/'; + if(data.ok) {{ + if(data.warning) alert(data.warning); + location.href = '/'; + }} else alert('Fehler: ' + data.error); }}) .catch(e => alert('Fehler: ' + e)); @@ -371,21 +378,28 @@ def page_search(db_path, query="", limit=30): """ if query: conn = _get_conn(db_path) - try: - rows = conn.execute( - "SELECT d.id, d.filename, d.file_type, d.word_count, " - "snippet(document_fts, 1, char(1), char(2), '...', 30) as snippet " - "FROM document_fts JOIN document_chunks dc ON document_fts.rowid = dc.id " - "JOIN documents d ON dc.doc_id = d.id " - "WHERE document_fts MATCH ? ORDER BY document_fts.rank LIMIT ?", - (query, limit) - ).fetchall() - except Exception: - like = f"%{query}%" + rows = None + # Erst FTS5-Syntax, bei Syntaxfehler (z.B. 'COVID-19', 'c++') + # quotierte Tokens, zuletzt LIKE-Fallback + for fts_query in (query, fts5_quote(query)): + try: + rows = conn.execute( + "SELECT d.id, d.filename, d.file_type, d.word_count, " + "snippet(document_fts, 1, char(1), char(2), '...', 30) as snippet " + "FROM document_fts JOIN document_chunks dc ON document_fts.rowid = dc.id " + "JOIN documents d ON dc.doc_id = d.id " + "WHERE document_fts MATCH ? ORDER BY document_fts.rank LIMIT ?", + (fts_query, limit) + ).fetchall() + break + except sqlite3.OperationalError: + continue + if rows is None: + like = f"%{escape_like(query)}%" rows = conn.execute( "SELECT DISTINCT d.id, d.filename, d.file_type, d.word_count, '' as snippet " "FROM documents d LEFT JOIN document_chunks dc ON dc.doc_id = d.id " - "WHERE dc.content LIKE ? OR d.filename LIKE ? LIMIT ?", + "WHERE dc.content LIKE ? ESCAPE '\\' OR d.filename LIKE ? ESCAPE '\\' LIMIT ?", (like, like, limit) ).fetchall() conn.close() @@ -566,31 +580,30 @@ def _reset_doc(self, doc_id): def _delete_file(self, doc_id): conn = _get_conn(self.db_path) try: - doc = conn.execute("SELECT file_path FROM documents WHERE id=?", (doc_id,)).fetchone() + doc = conn.execute("SELECT file_path, archived_path FROM documents WHERE id=?", + (doc_id,)).fetchone() if doc: - file_path = doc['file_path'] - trash_dir = Path(self.db_path).parent / "_Papierkorb" - trash_dir.mkdir(parents=True, exist_ok=True) - if os.path.exists(file_path): - base_name = os.path.basename(file_path) - trash_path = trash_dir / base_name - counter = 1 - while trash_path.exists(): - name, ext = os.path.splitext(base_name) - trash_path = trash_dir / f"{name}_{counter}{ext}" - counter += 1 + file_path = resolve_document_path(doc) + + # Erst DB-Eintrag in einer Transaktion entfernen, dann die + # Datei verschieben: schlaegt das Verschieben fehl, bleibt die + # DB trotzdem konsistent (vorher: Datei weg, DB-Eintrag da). + with conn: + conn.execute("DELETE FROM summaries WHERE source_type='document' AND source_id=?", (doc_id,)) + conn.execute("DELETE FROM document_keywords WHERE doc_id=?", (doc_id,)) + conn.execute("DELETE FROM document_chunks WHERE doc_id=?", (doc_id,)) + conn.execute("DELETE FROM digest_queue WHERE source_type='document' AND source_id=?", (doc_id,)) + conn.execute("DELETE FROM documents WHERE id=?", (doc_id,)) + + result = {"ok": True} + if file_path and os.path.exists(file_path): try: - shutil.move(file_path, str(trash_path)) + trash_path = move_to_trash(file_path, Path(self.db_path).parent / "_Papierkorb") + result["moved_to"] = str(trash_path) except Exception as e: - print(f"Move Error: {e}") - - conn.execute("DELETE FROM summaries WHERE source_type='document' AND source_id=?", (doc_id,)) - conn.execute("DELETE FROM document_keywords WHERE doc_id=?", (doc_id,)) - conn.execute("DELETE FROM document_chunks WHERE doc_id=?", (doc_id,)) - conn.execute("DELETE FROM digest_queue WHERE source_type='document' AND source_id=?", (doc_id,)) - conn.execute("DELETE FROM documents WHERE id=?", (doc_id,)) - conn.commit() - self._respond(200, json.dumps({"ok": True}), "application/json") + result["warning"] = (f"Aus der Datenbank entfernt, aber Datei konnte nicht " + f"in den Papierkorb verschoben werden: {e}") + self._respond(200, json.dumps(result), "application/json") else: self._respond(404, json.dumps({"ok": False, "error": "Doc not found"}), "application/json") except Exception as e: @@ -600,20 +613,16 @@ def _delete_file(self, doc_id): def _open_file(self, doc_id): conn = _get_conn(self.db_path) - row = conn.execute("SELECT file_path FROM documents WHERE id=?", (doc_id,)).fetchone() + row = conn.execute("SELECT file_path, archived_path FROM documents WHERE id=?", + (doc_id,)).fetchone() conn.close() if not row: return {"ok": False, "error": "Dokument nicht gefunden"} - file_path = row["file_path"] - p = Path(file_path) - if not p.exists(): + file_path = resolve_document_path(row) + if not file_path or not Path(file_path).exists(): return {"ok": False, "error": f"Datei nicht gefunden: {file_path}"} try: - if sys.platform == "win32": - os.startfile(str(p)) - else: - import subprocess - subprocess.run(["xdg-open", str(p)]) + open_path(file_path) return {"ok": True, "path": file_path} except Exception as e: return {"ok": False, "error": str(e)} diff --git a/wiki_indexer.py b/wiki_indexer.py index c3d27a0..3462e6b 100644 --- a/wiki_indexer.py +++ b/wiki_indexer.py @@ -18,15 +18,15 @@ import sqlite3 from pathlib import Path -from typing import Dict, Optional, Any, Union +from typing import Dict, Any, Union # Relative Imports (Paket-Kontext) mit Fallback auf absolute Imports (sys.path) try: - from .schema import ensure_schema + from .schema import ThreadLocalConnections from .chunker import chunk_text, estimate_tokens from .utils import sha256_hash as _sha256, extract_keywords as _extract_keywords except ImportError: - from schema import ensure_schema # type: ignore[no-redef] + from schema import ThreadLocalConnections # type: ignore[no-redef] from chunker import chunk_text, estimate_tokens # type: ignore[no-redef] from utils import sha256_hash as _sha256, extract_keywords as _extract_keywords # type: ignore[no-redef] @@ -42,19 +42,15 @@ class WikiIndexer: def __init__(self, knowledge_db: Path): self.knowledge_db = knowledge_db - self._conn: Optional[sqlite3.Connection] = None + self._conns = ThreadLocalConnections() # eine Connection pro Thread def _get_conn(self) -> sqlite3.Connection: """Lazy-init der DB-Connection mit Schema-Sicherstellung.""" - if self._conn is None: - self._conn = ensure_schema(self.knowledge_db) - return self._conn + return self._conns.get(self.knowledge_db) def close(self): """Schliesst DB-Connection.""" - if self._conn: - self._conn.close() - self._conn = None + self._conns.close_all() def index_from_bach(self, bach_db_path: Path, *, chunk_size: int = 350, diff --git a/zoll_station.py b/zoll_station.py index 249b853..bd536c0 100644 --- a/zoll_station.py +++ b/zoll_station.py @@ -34,6 +34,19 @@ def _require_script(path: Path) -> None: sys.exit(1) +def _digest_command(script_dir: Path, *args: str): + """Kommando + cwd fuer die digest-CLI. + + digest.py nutzt Paket-relative Imports (from .schema ...) und laeuft daher + nicht als Script (python digest.py), sondern nur als Modul: + python -m ... mit dem Elternordner als cwd. Liegt das Paket + in einem Ordner ohne gueltigen Modulnamen (z.B. ".db"), wird das per + pip installierte Paket "KnowledgeDigest" verwendet. + """ + package = script_dir.name if script_dir.name.isidentifier() else "KnowledgeDigest" + return [sys.executable, "-m", package, *args], str(script_dir.parent) + + def main(): parser = argparse.ArgumentParser(description="Wissensdatenbank Zoll-Station") parser.add_argument("--agent", required=True, choices=["claude", "gemini", "flash"], @@ -43,8 +56,7 @@ def main(): args = parser.parse_args() - script_dir = Path(__file__).parent - wissensdb_script = script_dir / "digest.py" + script_dir = Path(__file__).resolve().parent haiku_script = script_dir / "haiku_batch.py" # internal, gitignored # ASCII-only output: avoids UnicodeEncodeError on Windows consoles @@ -59,7 +71,8 @@ def main(): if os.environ.get("GEMINI_API_KEY"): print("Status: GEMINI_API_KEY gefunden. Initiierung des Hyper-Flash-Modes (API).") print("Fuehre API-Summarizer aus...\n") - subprocess.run([sys.executable, str(wissensdb_script), "summarize", "--flash", "--limit", str(limit)]) + cmd, cwd = _digest_command(script_dir, "summarize", "--flash", "--limit", str(limit)) + subprocess.run(cmd, cwd=cwd) else: print("Status: Kein GEMINI_API_KEY gefunden. Keine Abkuerzungen erlaubt. Manueller Uebersetzungs-Zoll faellig!\n") print("Generiere Schwarm-Prompt fuer manuelle Verarbeitung:\n")