Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion routes/personal_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Expand Down
17 changes: 15 additions & 2 deletions src/ai_interaction.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
# ---------------------------------------------------------------------------
Expand Down Expand Up @@ -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)"}
Expand Down
7 changes: 6 additions & 1 deletion src/chat_processor.py
Original file line number Diff line number Diff line change
Expand Up @@ -361,7 +361,12 @@ def build_context_preface(
# 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
Expand Down
38 changes: 30 additions & 8 deletions src/personal_docs.py
Original file line number Diff line number Diff line change
Expand Up @@ -220,6 +220,25 @@ def __init__(self, personal_dir: str, rag_manager=None):
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:
Expand Down Expand Up @@ -297,9 +316,10 @@ def add_directory(self, directory: str, *, index: bool = True, owner: str = None
# 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}")
Expand Down Expand Up @@ -328,9 +348,10 @@ def remove_directory(self, directory: str):
# 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:
Expand Down Expand Up @@ -417,7 +438,7 @@ def refresh_index(self):

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."""
Expand Down Expand Up @@ -447,7 +468,8 @@ def get_stats(self) -> Dict[str, Any]:

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

Expand All @@ -456,7 +478,7 @@ def index_all_directories(self):

# 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}")
Expand All @@ -472,7 +494,7 @@ def index_all_directories(self):
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}")
Expand Down
31 changes: 31 additions & 0 deletions tests/test_rag_live_manager_recovery.py
Original file line number Diff line number Diff line change
@@ -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