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