mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-08-09 10:39:11 +02:00
fix(rag): let runtime adopt the recovered document store
This commit is contained in:
parent
42da399b4d
commit
12da466572
5 changed files with 86 additions and 12 deletions
|
|
@ -143,7 +143,10 @@ def setup_personal_routes(personal_docs_manager, rag_manager, rag_available):
|
|||
|
||||
def _rag():
|
||||
"""Get the current RAG manager, retrying init if needed."""
|
||||
return get_rag_manager()
|
||||
recovered = get_rag_manager()
|
||||
if recovered is not None:
|
||||
personal_docs_manager.rag_manager = recovered
|
||||
return recovered
|
||||
|
||||
def _resolve_allowed_personal_dir(directory: str) -> str:
|
||||
"""Resolve a user-supplied personal-docs path under the allowed root."""
|
||||
|
|
|
|||
|
|
@ -67,6 +67,18 @@ def set_rag_manager(rag_mgr, personal_docs_mgr=None):
|
|||
_personal_docs_manager = personal_docs_mgr
|
||||
|
||||
|
||||
def _get_live_rag_manager():
|
||||
"""Resolve startup-degraded RAG through the shared personal-doc manager."""
|
||||
global _rag_manager
|
||||
|
||||
get_rag = getattr(_personal_docs_manager, "get_rag_manager", None)
|
||||
if callable(get_rag):
|
||||
recovered = get_rag()
|
||||
if recovered is not None:
|
||||
_rag_manager = recovered
|
||||
return _rag_manager
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Model resolution
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
@ -571,11 +583,12 @@ async def do_manage_rag(content: str, session_id: Optional[str] = None) -> Dict:
|
|||
if not os.path.isdir(directory):
|
||||
return {"error": f"Directory not found: {directory}"}
|
||||
|
||||
if not _rag_manager:
|
||||
rag_manager = _get_live_rag_manager()
|
||||
if not rag_manager:
|
||||
return {"error": "RAG manager not available"}
|
||||
|
||||
try:
|
||||
result = _rag_manager.index_personal_documents(directory)
|
||||
result = rag_manager.index_personal_documents(directory)
|
||||
indexed = result.get("indexed", 0) if isinstance(result, dict) else 0
|
||||
return {"action": "add_directory", "directory": directory,
|
||||
"results": f"Directory '{directory}' added to RAG index ({indexed} files indexed)"}
|
||||
|
|
|
|||
|
|
@ -361,7 +361,12 @@ class ChatProcessor:
|
|||
# RAG: search if enabled and rag_manager available, inject only above threshold
|
||||
if use_rag:
|
||||
try:
|
||||
rag_manager = getattr(self.personal_docs_manager, 'rag_manager', None)
|
||||
get_rag = getattr(self.personal_docs_manager, "get_rag_manager", None)
|
||||
rag_manager = (
|
||||
get_rag()
|
||||
if callable(get_rag)
|
||||
else getattr(self.personal_docs_manager, "rag_manager", None)
|
||||
)
|
||||
if rag_manager:
|
||||
results = rag_manager.search(message, k=5, owner=owner)
|
||||
# Filter by similarity threshold
|
||||
|
|
|
|||
|
|
@ -220,6 +220,25 @@ class PersonalDocsManager:
|
|||
self._load_excluded()
|
||||
self.refresh_index()
|
||||
|
||||
def get_rag_manager(self):
|
||||
"""Return the live RAG manager and adopt lazy startup recovery.
|
||||
|
||||
The application may start while Chroma is unavailable. Personal routes
|
||||
already retry the singleton later, but the long-lived manager retained
|
||||
the startup ``None`` and chat/agent retrieval stayed keyword-only. Keep
|
||||
the recovered singleton on this shared manager so every caller sees the
|
||||
same live provider.
|
||||
"""
|
||||
if self.rag_manager is not None:
|
||||
return self.rag_manager
|
||||
|
||||
from src.rag_singleton import get_rag_manager
|
||||
|
||||
recovered = get_rag_manager()
|
||||
if recovered is not None:
|
||||
self.rag_manager = recovered
|
||||
return recovered
|
||||
|
||||
def load_directories(self):
|
||||
"""Load the list of indexed directories from persistent storage."""
|
||||
try:
|
||||
|
|
@ -297,9 +316,10 @@ class PersonalDocsManager:
|
|||
# If RAG manager is available, index the directory immediately.
|
||||
# Callers that already indexed with owner metadata can pass
|
||||
# index=False so we do not create a second ownerless copy.
|
||||
if index and self.rag_manager:
|
||||
rag_manager = self.get_rag_manager() if index else None
|
||||
if rag_manager:
|
||||
try:
|
||||
result = self.rag_manager.index_personal_documents(directory, owner=owner)
|
||||
result = rag_manager.index_personal_documents(directory, owner=owner)
|
||||
logger.info(f"Indexed {result.get('indexed_count', 0)} chunks from {directory}")
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to index directory {directory}: {e}")
|
||||
|
|
@ -328,9 +348,10 @@ class PersonalDocsManager:
|
|||
# re-indexed only the remaining tracked dirs — ownerless and never
|
||||
# personal_dir — a catastrophic wipe (#1660). remove_directory now
|
||||
# removes exactly this directory's chunks and leaves the rest intact.
|
||||
if self.rag_manager:
|
||||
rag_manager = self.get_rag_manager()
|
||||
if rag_manager:
|
||||
try:
|
||||
self.rag_manager.remove_directory(directory)
|
||||
rag_manager.remove_directory(directory)
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to remove directory from RAG index: {e}")
|
||||
else:
|
||||
|
|
@ -417,7 +438,7 @@ class PersonalDocsManager:
|
|||
|
||||
def retrieve(self, query: str, k: int = 5) -> List[str]:
|
||||
"""Retrieve relevant documents for a query."""
|
||||
return retrieve_personal(self.index, query, k, self.rag_manager)
|
||||
return retrieve_personal(self.index, query, k, self.get_rag_manager())
|
||||
|
||||
def get_file_list(self) -> List[Dict[str, Any]]:
|
||||
"""Get list of indexed files with metadata."""
|
||||
|
|
@ -447,7 +468,8 @@ class PersonalDocsManager:
|
|||
|
||||
def index_all_directories(self):
|
||||
"""Re-index all tracked directories in the RAG system."""
|
||||
if not self.rag_manager:
|
||||
rag_manager = self.get_rag_manager()
|
||||
if not rag_manager:
|
||||
logger.warning("No RAG manager available for indexing")
|
||||
return
|
||||
|
||||
|
|
@ -456,7 +478,7 @@ class PersonalDocsManager:
|
|||
|
||||
# Index the base personal directory
|
||||
try:
|
||||
result = self.rag_manager.index_personal_documents(self.personal_dir)
|
||||
result = rag_manager.index_personal_documents(self.personal_dir)
|
||||
if result.get('success'):
|
||||
success_count += 1
|
||||
logger.info(f"Indexed base directory: {self.personal_dir}")
|
||||
|
|
@ -472,7 +494,7 @@ class PersonalDocsManager:
|
|||
continue
|
||||
|
||||
try:
|
||||
result = self.rag_manager.index_personal_documents(directory)
|
||||
result = rag_manager.index_personal_documents(directory)
|
||||
if result.get('success'):
|
||||
success_count += 1
|
||||
logger.info(f"Indexed directory: {directory}")
|
||||
|
|
|
|||
31
tests/test_rag_live_manager_recovery.py
Normal file
31
tests/test_rag_live_manager_recovery.py
Normal file
|
|
@ -0,0 +1,31 @@
|
|||
from src import ai_interaction
|
||||
from src import personal_docs
|
||||
|
||||
|
||||
class _RecoveredRag:
|
||||
def search(self, query, k=5, owner=None):
|
||||
return []
|
||||
|
||||
|
||||
def test_personal_docs_manager_adopts_recovered_singleton(monkeypatch):
|
||||
recovered = _RecoveredRag()
|
||||
manager = personal_docs.PersonalDocsManager.__new__(personal_docs.PersonalDocsManager)
|
||||
manager.rag_manager = None
|
||||
monkeypatch.setattr("src.rag_singleton.get_rag_manager", lambda: recovered)
|
||||
|
||||
assert manager.get_rag_manager() is recovered
|
||||
assert manager.rag_manager is recovered
|
||||
|
||||
|
||||
def test_ai_interaction_uses_shared_recovered_manager(monkeypatch):
|
||||
recovered = _RecoveredRag()
|
||||
personal_manager = type(
|
||||
"PersonalManager",
|
||||
(),
|
||||
{"get_rag_manager": lambda self: recovered},
|
||||
)()
|
||||
monkeypatch.setattr(ai_interaction, "_rag_manager", None)
|
||||
monkeypatch.setattr(ai_interaction, "_personal_docs_manager", personal_manager)
|
||||
|
||||
assert ai_interaction._get_live_rag_manager() is recovered
|
||||
assert ai_interaction._rag_manager is recovered
|
||||
Loading…
Add table
Add a link
Reference in a new issue