Initial commit: library-desk service extraction from portainer-core
Build and Push / build (release) Successful in 36s

This commit is contained in:
2025-12-11 17:28:23 +01:00
commit 95852190ba
59 changed files with 17146 additions and 0 deletions
View File
+932
View File
@@ -0,0 +1,932 @@
"""
Knowledge Consolidation Service (Librarian Logic)
Processes unprocessed SearchQuery nodes from HybridRAG searches
to consolidate new knowledge into wiki pages.
This service:
1. Queries Neo4j for unprocessed SearchQuery nodes
2. Analyzes web results with Ollama for novel information
3. Creates/updates wiki pages with new facts
4. Updates knowledge graph with new entities
5. Marks SearchQuery nodes as processed
"""
import logging
import json
from datetime import datetime, timedelta
from typing import List, Dict, Any, Optional
from src.clients.neo4j_client import Neo4jClient
from src.clients.ollama_client import OllamaClient
from src.clients.wikijs_client import WikiJSClient
from src.services.wiki_page_writer import WikiPageWriter
from src.models.consolidation import (
SearchQueryInfo,
ConsolidationResult,
ConsolidationResponse
)
from src.config import Settings
logger = logging.getLogger(__name__)
class ConsolidationService:
"""
Service for consolidating knowledge from search results.
"""
def __init__(
self,
neo4j: Neo4jClient,
ollama: OllamaClient,
wiki: WikiJSClient,
settings: Settings,
ingestion_service: Optional["IngestionService"] = None
):
self.neo4j = neo4j
self.ollama = ollama
self.wiki = wiki
self.settings = settings
self.wiki_page_writer = WikiPageWriter(ollama_client=ollama)
self.ingestion_service = ingestion_service # Optional to avoid circular dependency
async def consolidate_knowledge(
self,
process_limit: int = 10,
lookback_days: int = 7,
min_web_results: int = 2,
dry_run: bool = False
) -> ConsolidationResponse:
"""
Process unprocessed search queries and consolidate knowledge.
Args:
process_limit: Maximum searches to process
lookback_days: Only process searches from last N days
min_web_results: Minimum web results required to consolidate
dry_run: If True, analyze but don't create pages
Returns:
ConsolidationResponse with processing results
"""
logger.info(f"Starting knowledge consolidation")
logger.info(f"Limits: process={process_limit}, lookback={lookback_days}d, min_web={min_web_results}")
if dry_run:
logger.warning("DRY RUN MODE - will not create wiki pages")
# Find unprocessed searches
unprocessed = await self._find_unprocessed_searches(lookback_days, process_limit)
if not unprocessed:
logger.info("No unprocessed searches found")
return ConsolidationResponse(
total_found=0,
processed_count=0,
pages_created=0,
pages_updated=0,
entities_added=0,
errors=[],
results=[],
dry_run=dry_run
)
logger.info(f"Found {len(unprocessed)} unprocessed searches")
# Process each search
results: List[ConsolidationResult] = []
total_pages_created = 0
total_pages_updated = 0
total_entities_added = 0
errors: List[str] = []
for search in unprocessed:
try:
result = await self._process_search(
search=search,
min_web_results=min_web_results,
dry_run=dry_run
)
if result:
results.append(result)
total_pages_created += result.pages_created
total_pages_updated += result.pages_updated
total_entities_added += result.entities_added
# Mark as processed if not dry run (even if skipped)
# This prevents searches from accumulating when they don't meet criteria
if not dry_run:
await self._mark_search_processed(search['id'])
except Exception as e:
error_msg = f"Search {search['id'][:8]}: {str(e)}"
logger.error(f"Failed to process search: {error_msg}", exc_info=True)
errors.append(error_msg)
results.append(ConsolidationResult(
search_id=search['id'],
query=search['query'],
error=str(e)
))
# Mark as processed even on error (to avoid retrying failed searches forever)
if not dry_run:
await self._mark_search_processed(search['id'])
# Build response
processed_count = len([r for r in results if not r.error])
response = ConsolidationResponse(
total_found=len(unprocessed),
processed_count=processed_count,
pages_created=total_pages_created,
pages_updated=total_pages_updated,
entities_added=total_entities_added,
errors=errors,
results=results,
dry_run=dry_run
)
logger.info(
f"Consolidation complete: {processed_count}/{len(unprocessed)} searches, "
f"{total_pages_created} pages created, {total_pages_updated} updated, "
f"{total_entities_added} entities added"
)
return response
async def _find_unprocessed_searches(
self,
lookback_days: int,
limit: int
) -> List[Dict[str, Any]]:
"""
Find unprocessed SearchQuery nodes from Neo4j.
"""
lookback_date = datetime.now() - timedelta(days=lookback_days)
query = """
MATCH (sq:SearchQuery {processed: false})
WHERE sq.timestamp > datetime($lookback_date)
RETURN sq.id as id,
sq.query as query,
sq.user as user,
sq.timestamp as timestamp,
sq.total_results as total_results,
sq.web_count as web_count,
sq.keywords as keywords
ORDER BY sq.timestamp DESC
LIMIT $limit
"""
try:
results = await self.neo4j.execute_query(
query,
{
"lookback_date": lookback_date.isoformat(),
"limit": limit
}
)
searches = []
for record in results:
searches.append({
'id': record['id'],
'query': record['query'],
'user': record['user'],
'timestamp': record['timestamp'],
'total_results': record.get('total_results', 0),
'web_count': record.get('web_count', 0),
'keywords': record.get('keywords', [])
})
return searches
except Exception as e:
logger.error(f"Failed to find unprocessed searches: {e}")
return []
async def _process_search(
self,
search: Dict[str, Any],
min_web_results: int,
dry_run: bool
) -> Optional[ConsolidationResult]:
"""
Process a single search query for knowledge consolidation.
"""
search_id = search['id']
query = search['query']
user = search['user']
web_count = search.get('web_count', 0)
logger.info(f"Processing: '{query}' (user: {user}, web: {web_count})")
# Skip if insufficient web results
if web_count < min_web_results:
logger.info(f"Skipping - insufficient web results ({web_count} < {min_web_results})")
return None
# Get web results from SearchQuery
web_results = await self._get_web_results(search_id)
if not web_results:
logger.info("No web results found in database")
return None
logger.info(f"Retrieved {len(web_results)} web results")
# Analyze web results with Ollama for novel information
analysis = await self._analyze_web_results(
query=query,
web_results=web_results,
keywords=search.get('keywords', []),
user=user
)
if not analysis or not analysis.get('has_novel_info'):
logger.info("No novel information found")
return ConsolidationResult(
search_id=search_id,
query=query
)
# Extract consolidation actions
pages_to_create = analysis.get('new_pages', [])
pages_to_update = analysis.get('update_pages', [])
new_entities = analysis.get('new_entities', [])
logger.info(
f"Analysis: {len(pages_to_create)} new pages, "
f"{len(pages_to_update)} updates, {len(new_entities)} entities"
)
if dry_run:
logger.info("[DRY RUN] Would create/update pages and entities")
return ConsolidationResult(
search_id=search_id,
query=query,
pages_created=len(pages_to_create),
pages_updated=len(pages_to_update),
entities_added=len(new_entities)
)
# Create/update wiki pages
pages_created = 0
pages_updated = 0
entities_added = 0
# Create new pages
for page_data in pages_to_create:
try:
await self._create_or_consolidate_page(
user=user,
title=page_data.get('title'),
path=page_data.get('path'),
summary=page_data.get('summary'),
source_query=query,
web_results=web_results
)
pages_created += 1
logger.info(f"Created page: {page_data.get('title')}")
except Exception as e:
logger.error(f"Failed to create page {page_data.get('title')}: {e}")
# Update existing pages
for page_data in pages_to_update:
try:
await self._update_page_with_facts(
title=page_data.get('title'),
new_facts=page_data.get('new_facts', []),
source_url=page_data.get('source_url'),
user=user
)
pages_updated += 1
logger.info(f"Updated page: {page_data.get('title')}")
except Exception as e:
logger.error(f"Failed to update page {page_data.get('title')}: {e}")
# Add new entities to graph
for entity_data in new_entities:
try:
await self._add_entity_to_graph(
user=user,
entity_name=entity_data.get('name'),
entity_type=entity_data.get('type'),
description=entity_data.get('description'),
source_search_id=search_id
)
entities_added += 1
logger.info(f"Added entity: {entity_data.get('name')}")
except Exception as e:
logger.error(f"Failed to add entity {entity_data.get('name')}: {e}")
return ConsolidationResult(
search_id=search_id,
query=query,
pages_created=pages_created,
pages_updated=pages_updated,
entities_added=entities_added
)
async def _get_web_results(self, search_id: str) -> List[Dict[str, Any]]:
"""Get web results for a search from Neo4j."""
query = """
MATCH (sq:SearchQuery {id: $search_id})-[f:FOUND]->(wr:WebResult)
RETURN wr.url as url,
wr.title as title,
wr.content as content,
f.rank as rank,
f.rrf_score as rrf_score
ORDER BY f.rank
LIMIT 20
"""
try:
results = await self.neo4j.execute_query(query, {"search_id": search_id})
web_results = []
for record in results:
web_results.append({
'url': record['url'],
'title': record['title'],
'content': record['content'],
'rank': record['rank'],
'rrf_score': record['rrf_score']
})
return web_results
except Exception as e:
logger.error(f"Failed to get web results: {e}")
return []
async def _analyze_web_results(
self,
query: str,
web_results: List[Dict[str, Any]],
keywords: List[str],
user: str = "jpmschweitzer"
) -> Optional[Dict[str, Any]]:
"""
Analyze web results with Ollama for novel information.
Returns analysis with has_novel_info, new_pages, update_pages, new_entities.
"""
# Fetch existing taxonomy structure for this user
try:
taxonomy_structure = await self.wiki.get_taxonomy_structure(f"users/{user}")
existing_paths_info = self._format_taxonomy_for_prompt(taxonomy_structure)
logger.info(f"Fetched taxonomy with {len(taxonomy_structure)} categories for user {user}")
except Exception as e:
logger.warning(f"Failed to fetch taxonomy structure: {e}")
existing_paths_info = ""
# Build analysis prompt
web_summary = "\n\n".join([
f"[{i+1}] {r['title']}\n{r['url']}\n{r['content'][:300]}..."
for i, r in enumerate(web_results[:5])
])
prompt = f"""You are a Librarian helping build a personal knowledge base and extended memory system.
Analyze these web search results for information worth documenting in our personal wiki.
Query: "{query}"
Keywords: {', '.join(keywords) if keywords else 'none'}
Web Results:
{web_summary}
This is a PERSONAL knowledge base using Schema.org-aligned taxonomy that captures:
- People: Family members, friends, colleagues, public figures (Schema.org: Person)
- Companies: Businesses, organizations, institutions (Schema.org: Organization)
- Places: Locations, restaurants, travel destinations (Schema.org: Place)
- Entertainment: Books, movies, TV, music, games (Schema.org: CreativeWork)
- Recipes: Food, cooking techniques, ingredients (Schema.org: CreativeWork/Recipe)
- Products: Purchased items, gear, tools, equipment (Schema.org: Product)
- Technology: Software, applications, infrastructure (Schema.org: SoftwareApplication)
- Health: Medical info, fitness, wellness (Schema.org: MedicalEntity)
- Events: Concerts, travel, appointments, important dates (Schema.org: Event)
- Hobbies: Personal interests, activities, pastimes (Custom extension)
- Projects: Work projects, personal projects (Schema.org: Project)
- Reference: General knowledge, how-tos (Custom extension)
Identify information worth documenting:
1. New topics/people/things that deserve their own wiki page
2. Facts that could enhance existing pages
3. Entities (people, places, things, concepts) for the knowledge graph
Be INCLUSIVE - if someone searched for it, it's likely worth documenting.
Personal information is just as valuable as technical information.
**CRITICAL: Use ONLY these Schema.org-aligned path prefixes (case-sensitive):**
- People: `people/<name>` (Schema.org: Person)
- Companies: `companies/<company-name>` (Schema.org: Organization)
- Places: `places/<location>` (Schema.org: Place)
- Entertainment (Schema.org: CreativeWork):
- Books: `entertainment/books/<title>`
- Movies: `entertainment/movies/<title>`
- TV: `entertainment/tv/<title>`
- Music: `entertainment/music/<artist-or-album>`
- Games: `entertainment/games/<title>`
- Recipes: `recipes/<cuisine-or-category>/<dish>` (Schema.org: Recipe)
- Products: `products/<category>/<product-name>` (Schema.org: Product)
- Technology: `technology/<category>/<topic>` (Schema.org: SoftwareApplication)
- Health: `health/<category>/<topic>` (Schema.org: MedicalEntity)
- Events: `events/<event-type>/<event-name>` (Schema.org: Event)
- Hobbies: `hobbies/<hobby-name>` (Custom extension)
- Projects: `projects/<project-name>` (Schema.org: Project)
- Reference: `reference/<category>/<topic>` (Custom extension)
**Path Rules:**
- Use lowercase with hyphens (kebab-case): "machine-learning" not "Machine_Learning"
- Keep paths 2-3 levels deep maximum
- Be consistent with existing paths when possible
{existing_paths_info}
Return ONLY valid JSON:
{{
"has_novel_info": true,
"new_pages": [
{{"title": "Page Title", "path": "companies/example-company", "summary": "What information to include"}}
],
"update_pages": [
{{"title": "Existing Page", "new_facts": ["fact 1"], "source_url": "url"}}
],
"new_entities": [
{{"name": "Entity Name", "type": "person/place/thing/concept/recipe/media", "description": "Brief description"}}
]
}}
JSON:"""
try:
# Call Ollama for analysis
response = await self.ollama.generate_text(
prompt=prompt,
model=self.settings.reranker_model, # Use mistral-nemo
stream=False
)
if not response:
logger.warning("Empty response from Ollama")
return None
# Extract JSON from response
response_clean = response.strip()
if '{' in response_clean:
json_start = response_clean.find('{')
json_end = response_clean.rfind('}') + 1
response_clean = response_clean[json_start:json_end]
analysis = json.loads(response_clean)
return analysis
except json.JSONDecodeError as e:
logger.error(f"Failed to parse Ollama response as JSON: {e}")
logger.debug(f"Response was: {response[:500]}")
return None
except Exception as e:
logger.error(f"Analysis failed: {e}", exc_info=True)
return None
def _format_taxonomy_for_prompt(self, taxonomy: Dict[str, List[str]]) -> str:
"""
Format taxonomy structure for inclusion in LLM prompt.
Args:
taxonomy: Dict mapping categories to subcategories
Returns:
Formatted string showing existing paths
"""
if not taxonomy:
return ""
lines = ["**Existing paths in your wiki (PREFER these over creating new ones):**"]
for category, subcategories in taxonomy.items():
if subcategories:
lines.append(f"- {category}/")
for sub in subcategories:
lines.append(f" - {category}/{sub}/")
else:
lines.append(f"- {category}/")
lines.append("")
lines.append("**IMPORTANT:** If a suitable existing path exists, use it instead of creating a new category.")
lines.append("Example: NATO should go in `reference/political-entities/` not a new `reference/military-alliances/`")
return "\n".join(lines)
async def _mark_search_processed(self, search_id: str):
"""Mark SearchQuery node as processed."""
query = """
MATCH (sq:SearchQuery {id: $search_id})
SET sq.processed = true,
sq.processed_at = datetime()
RETURN sq.id
"""
try:
await self.neo4j.execute_query(query, {"search_id": search_id})
logger.debug(f"Marked search {search_id} as processed")
except Exception as e:
logger.error(f"Failed to mark search as processed: {e}")
async def _apply_bidirectional_entity_linking(
self,
page_id: int,
page_title: str,
user: str
) -> Dict[str, int]:
"""
Apply bidirectional entity linking after page creation/update.
This runs AFTER ingestion so entities are extracted and in the graph.
Steps:
1. Link entities in the new page (forward links to existing entities)
2. Find pages that mention the new entity (reverse references)
3. Link entities in those pages (backward links to the new entity)
Args:
page_id: Wiki page ID
page_title: Page title (used to find reverse references)
user: User identifier
Returns:
Dict with link counts: {
"forward_links": int, # Links added to the new page
"backward_links": int, # Links added to other pages pointing to new page
"pages_updated": int # Number of other pages updated
}
"""
from src.core.multi_tenancy import get_neo4j_user_base_label
forward_links = 0
backward_links = 0
pages_updated = 0
try:
# Import here to avoid circular dependency
from src.routers.entity_linking import link_entities_in_page, EntityLinkingRequest
from src.core.dependencies import get_wiki_service, get_graph_service
wiki_service = get_wiki_service()
graph_service = get_graph_service()
# STEP 1: Forward linking - link entities in the new page
logger.info(f"Step 1/3: Linking entities in page {page_id} ('{page_title}')")
try:
forward_result = await link_entities_in_page(
request=EntityLinkingRequest(
user=user,
page_id=page_id,
create_relationships=True,
re_index_if_changed=False # Already indexed, no need to re-index
),
wiki_service=wiki_service,
graph_service=graph_service,
ingestion_service=self.ingestion_service,
api_key="" # Internal call, no auth needed
)
forward_links = forward_result.content_links_added
logger.info(f"Added {forward_links} forward links in page {page_id}")
except Exception as e:
logger.error(f"Failed to add forward links: {e}")
# STEP 2: Find reverse references - which pages mention this new entity?
logger.info(f"Step 2/3: Finding pages that mention '{page_title}'")
user_base_label = get_neo4j_user_base_label(user)
# Query to find documents that mention entities with this page's title
reverse_query = f"""
// Find entities with the same name as the page title
MATCH (e:{user_base_label})
WHERE toLower(e.name) = toLower($title)
AND NOT e:Document
// Find documents that mention those entities
MATCH (d:Document)-[r:MENTIONS]->(e)
WHERE d.page_id <> $page_id // Exclude the page itself
RETURN DISTINCT d.page_id as page_id, d.title as title
LIMIT 50
"""
try:
reverse_refs = await self.neo4j.execute_query(
reverse_query,
{"title": page_title, "page_id": page_id}
)
logger.info(f"Found {len(reverse_refs)} pages that mention '{page_title}'")
except Exception as e:
logger.error(f"Failed to find reverse references: {e}")
reverse_refs = []
# STEP 3: Backward linking - add links in those pages to the new entity
if reverse_refs:
logger.info(f"Step 3/3: Adding backward links in {len(reverse_refs)} pages")
for ref in reverse_refs:
try:
backward_result = await link_entities_in_page(
request=EntityLinkingRequest(
user=user,
page_id=ref['page_id'],
create_relationships=False, # Relationships already exist
re_index_if_changed=False # Don't re-index for link updates
),
wiki_service=wiki_service,
graph_service=graph_service,
ingestion_service=self.ingestion_service,
api_key=""
)
if backward_result.content_links_added > 0:
backward_links += backward_result.content_links_added
pages_updated += 1
logger.info(
f"Added {backward_result.content_links_added} links "
f"in page {ref['page_id']} ('{ref['title']}')"
)
except Exception as e:
logger.error(f"Failed to add backward links in page {ref['page_id']}: {e}")
else:
logger.info("Step 3/3: No reverse references found, skipping backward linking")
return {
"forward_links": forward_links,
"backward_links": backward_links,
"pages_updated": pages_updated
}
except Exception as e:
logger.error(f"Bidirectional entity linking failed: {e}", exc_info=True)
return {
"forward_links": 0,
"backward_links": 0,
"pages_updated": 0
}
async def _create_or_consolidate_page(
self,
user: str,
title: str,
path: str,
summary: str,
source_query: str,
web_results: List[Dict[str, Any]]
):
"""
Create wiki page or consolidate with existing synonym page.
Uses WikiPageWriter for intelligent LLM-based content generation:
- For new pages: Holistic structured content creation
- For existing pages: Zero-loss reconstruction with conflict detection
"""
# Normalize path to user namespace
if not path.startswith(f"users/{user}"):
path = f"users/{user}/{path.lstrip('/')}"
# Format web results as source information
source_information = [
{
'title': r['title'],
'url': r['url'],
'content': r['content']
}
for r in web_results[:5] # Top 5 web results
]
# Search for existing pages with similar titles (synonym consolidation)
existing_pages = await self.wiki.search_pages(title, path_prefix=f"users/{user}")
if existing_pages:
# Page exists - reconstruct with new information using LLM
logger.info(f"Found existing page for '{title}', will reconstruct with new info")
page_id = existing_pages[0]['id']
# Get current content
existing_page = await self.wiki.get_page(page_id)
if existing_page:
# Build new information text from summary and web results
new_information = f"{summary}\n\n"
for r in web_results[:3]:
new_information += f"- {r['title']}: {r['content'][:200]}...\n"
# Use WikiPageWriter to reconstruct with LLM
reconstructed_content, conflicts = await self.wiki_page_writer.reconstruct_page(
title=title,
existing_content=existing_page['content'],
new_information=new_information,
new_sources=source_information,
detect_conflicts=True
)
if conflicts:
logger.warning(
f"Detected {len(conflicts)} conflicts when updating '{title}' - "
"LLM chose most authoritative sources"
)
await self.wiki.update_page(
page_id=page_id,
content=reconstructed_content
)
logger.info(f"Reconstructed existing page: {title}")
# Trigger ingestion to update vectors and graph
if self.ingestion_service:
try:
await self.ingestion_service.ingest_page(
page_id=page_id,
user=user,
force_refresh=True
)
logger.info(f"Ingested updated page {page_id} into knowledge base")
# Apply bidirectional entity linking after ingestion
link_stats = await self._apply_bidirectional_entity_linking(
page_id=page_id,
page_title=title,
user=user
)
logger.info(
f"Entity linking complete: {link_stats['forward_links']} forward links, "
f"{link_stats['backward_links']} backward links "
f"({link_stats['pages_updated']} pages updated)"
)
except Exception as e:
logger.error(f"Failed to ingest updated page {page_id}: {e}")
return
# Create new page with LLM-generated structured content
logger.info(f"Creating new page: {title}")
# Use WikiPageWriter to create structured content
content = await self.wiki_page_writer.create_page(
title=title,
topic_summary=summary,
source_information=source_information,
entities=None, # Could extract from keywords if available
related_docs=None
)
# Extract tags from path for dossier organization
path_parts = path.split('/')
tags = [part for part in path_parts if part and part not in ['users', user]]
created_page = await self.wiki.create_page(
path=path,
title=title,
content=content,
description=f"Consolidated from search: {source_query}",
tags=tags[:3], # Limit to 3 tags
is_published=True
)
page_id = created_page.get("id") if created_page else None
logger.info(f"Created new page: {path} (page_id: {page_id})")
# Trigger ingestion to update vectors and graph
if self.ingestion_service and page_id:
try:
await self.ingestion_service.ingest_page(
page_id=page_id,
user=user,
force_refresh=False # New page, no need to force
)
logger.info(f"Ingested new page {page_id} into knowledge base")
# Apply bidirectional entity linking after ingestion
link_stats = await self._apply_bidirectional_entity_linking(
page_id=page_id,
page_title=title,
user=user
)
logger.info(
f"Entity linking complete: {link_stats['forward_links']} forward links, "
f"{link_stats['backward_links']} backward links "
f"({link_stats['pages_updated']} pages updated)"
)
except Exception as e:
logger.error(f"Failed to ingest new page {page_id}: {e}")
async def _update_page_with_facts(
self,
title: str,
new_facts: List[str],
source_url: str,
user: str
):
"""
Update existing page with new facts using LLM reconstruction.
Uses WikiPageWriter to intelligently merge facts with zero loss.
"""
# Search for page
pages = await self.wiki.search_pages(title, path_prefix=f"users/{user}")
if not pages:
logger.warning(f"Page '{title}' not found for update")
return
page_id = pages[0]['id']
existing_page = await self.wiki.get_page(page_id)
if not existing_page:
return
# Build new information from facts
new_information = "\n".join([f"- {fact}" for fact in new_facts])
# Format source
source_information = [{
'title': source_url,
'url': source_url,
'content': new_information
}]
# Use WikiPageWriter to reconstruct with LLM
reconstructed_content, conflicts = await self.wiki_page_writer.reconstruct_page(
title=title,
existing_content=existing_page['content'],
new_information=new_information,
new_sources=source_information,
detect_conflicts=True
)
if conflicts:
logger.warning(
f"Detected {len(conflicts)} conflicts when updating '{title}' with new facts"
)
await self.wiki.update_page(
page_id=page_id,
content=reconstructed_content
)
# Trigger ingestion to update vectors and graph
logger.debug(f"ingestion_service available: {self.ingestion_service is not None}")
if self.ingestion_service:
try:
logger.info(f"Starting ingestion for updated page {page_id}")
await self.ingestion_service.ingest_page(
page_id=page_id,
user=user,
force_refresh=True
)
logger.info(f"Ingested updated page {page_id} into knowledge base")
# Apply bidirectional entity linking after ingestion
link_stats = await self._apply_bidirectional_entity_linking(
page_id=page_id,
page_title=title,
user=user
)
logger.info(
f"Entity linking complete: {link_stats['forward_links']} forward links, "
f"{link_stats['backward_links']} backward links "
f"({link_stats['pages_updated']} pages updated)"
)
except Exception as e:
logger.error(f"Failed to ingest updated page {page_id}: {e}")
async def _add_entity_to_graph(
self,
user: str,
entity_name: str,
entity_type: str,
description: str,
source_search_id: str
):
"""Add new entity to knowledge graph."""
from src.core.multi_tenancy import get_neo4j_user_base_label
user_base_label = get_neo4j_user_base_label(user)
# Create entity node with appropriate type label
type_label = entity_type.capitalize() if entity_type else "Entity"
query = f"""
MERGE (e:{user_base_label}:{type_label} {{name: $name}})
ON CREATE SET
e.description = $description,
e.created_at = datetime(),
e.source = 'librarian_consolidation',
e.source_search_id = $search_id
ON MATCH SET
e.updated_at = datetime()
RETURN e
"""
try:
await self.neo4j.execute_query(query, {
"name": entity_name,
"description": description,
"search_id": source_search_id
})
logger.debug(f"Added entity to graph: {entity_name} ({entity_type})")
except Exception as e:
logger.error(f"Failed to add entity to graph: {e}")
File diff suppressed because it is too large Load Diff
+739
View File
@@ -0,0 +1,739 @@
"""
HybridRAG service combining vector, graph, and web search.
6-Phase Pipeline:
0. Query Enhancement - Extract keywords/synonyms with LLM
1. Parallel Retrieval - Vector + Graph + Web search
2. RRF Fusion - Merge results with Reciprocal Rank Fusion
3. Enrichment - Add related dossiers via graph
4. LLM Re-ranking - Re-rank with mistral-nemo
5. Context Formatting - Format for LLM consumption
6. Persistence - Store for Librarian processing
"""
import asyncio
import time
import json
import uuid
from typing import List, Dict, Any, Optional
import logging
from src.services.vector_service import VectorService
from src.services.graph_service import GraphService
from src.clients.searxng_client import SearXNGClient
from src.clients.ollama_client import OllamaClient
from src.config import Settings
from src.models.hybrid_rag import (
HybridRAGConfig, HybridRAGRequest, HybridRAGResponse,
HybridRAGResult, TimingBreakdown, KeywordExtraction,
RelatedDossier
)
from src.core.multi_tenancy import get_neo4j_user_base_label, get_neo4j_user_label
logger = logging.getLogger(__name__)
class HybridRAGService:
"""
Service for HybridRAG multi-source search with fusion and re-ranking.
"""
def __init__(
self,
vector_service: VectorService,
graph_service: GraphService,
searxng_client: SearXNGClient,
ollama_client: OllamaClient,
settings: Settings
):
"""
Initialize HybridRAG service.
Args:
vector_service: Service for Qdrant vector search
graph_service: Service for Neo4j graph search
searxng_client: Client for web search
ollama_client: Client for LLM (keyword extraction, re-ranking)
settings: Application settings
"""
self.vector = vector_service
self.graph = graph_service
self.searxng = searxng_client
self.ollama = ollama_client
self.settings = settings
self.reranker_model = settings.reranker_model
async def search(
self,
query: str,
user: str,
config: Optional[HybridRAGConfig] = None
) -> HybridRAGResponse:
"""
Execute HybridRAG search across all sources.
Args:
query: Search query
user: User identifier
config: Optional configuration override
Returns:
Complete search response with ranked results and timing
"""
start_time = time.time()
timing = {}
# Use default config if not provided
if not config:
config = HybridRAGConfig()
logger.info(f"HybridRAG search: '{query}' for user '{user}'")
# Phase 0: Query Enhancement
phase0_start = time.time()
keywords_data = await self._extract_keywords_and_synonyms(query)
timing["query_enhancement_ms"] = (time.time() - phase0_start) * 1000
# Phase 1: Parallel Retrieval
phase1_start = time.time()
raw_results = await self._retrieve_parallel(query, user, config, keywords_data)
timing["vector_ms"] = raw_results.get("timing", {}).get("vector_ms", 0)
timing["graph_ms"] = raw_results.get("timing", {}).get("graph_ms", 0)
timing["web_ms"] = raw_results.get("timing", {}).get("web_ms", 0)
# Phase 2: RRF Fusion
phase2_start = time.time()
fused_results = self._reciprocal_rank_fusion(
results_by_source={
"vector": raw_results.get("vector", []),
"graph": raw_results.get("graph", []),
"web": raw_results.get("web", [])
},
k=config.rrf_k
)
timing["fusion_ms"] = (time.time() - phase2_start) * 1000
# Phase 3: Enrichment
phase3_start = time.time()
if config.enable_enrichment:
enriched_results = await self._enrich_with_related_dossiers(fused_results, user)
else:
enriched_results = fused_results
timing["enrichment_ms"] = (time.time() - phase3_start) * 1000
# Phase 4: LLM Re-ranking
phase4_start = time.time()
if config.enable_reranking and len(enriched_results) > 1:
reranked_results = await self._rerank_with_llm(enriched_results[:20], query)
else:
reranked_results = enriched_results
timing["reranking_ms"] = (time.time() - phase4_start) * 1000
# Limit to final result count
final_results = reranked_results[:config.final_result_count]
# Update final ranks
for i, result in enumerate(final_results, start=1):
result["final_rank"] = i
# Convert to HybridRAGResult models
result_models = self._convert_to_result_models(final_results)
# Phase 5: Context Formatting
context = self._format_context_for_llm(result_models)
# Calculate source counts
source_counts = {}
for result in result_models:
for source in result.sources:
source_counts[source] = source_counts.get(source, 0) + 1
timing["total_ms"] = (time.time() - start_time) * 1000
# Phase 6: Persistence (async, non-blocking)
phase6_start = time.time()
search_id = await self._persist_search_for_librarian(
query=query,
user=user,
keywords_data=keywords_data,
raw_results=raw_results,
final_results=final_results,
timing=timing
)
timing["persistence_ms"] = (time.time() - phase6_start) * 1000
# Build response
return HybridRAGResponse(
query=query,
keywords=KeywordExtraction(**keywords_data),
results=result_models,
context=context,
source_counts=source_counts,
total_results=len(result_models),
timing=TimingBreakdown(**timing),
config_used=config,
search_id=search_id
)
async def _extract_keywords_and_synonyms(self, query: str) -> Dict[str, Any]:
"""
Phase 0: Extract keywords, entities, and synonyms using LLM.
Args:
query: Search query
Returns:
Dictionary with keywords, entities, synonyms, expansions
"""
prompt = f"""Extract search terms from this query. For each important word, provide synonyms and expansions.
Query: "{query}"
Return ONLY valid JSON:
{{
"core_keywords": ["key", "words", "from", "query"],
"synonyms": {{
"word": ["alternative", "terms"]
}}
}}
Example for "Docker container hosting":
{{
"core_keywords": ["docker", "container", "hosting"],
"synonyms": {{
"docker": ["containerization", "container runtime"],
"hosting": ["server", "infrastructure"]
}}
}}
JSON:"""
try:
response = await self.ollama.generate_text(
prompt=prompt,
model=self.reranker_model
)
# Parse JSON response (handle potential extra text)
response_clean = response.strip()
# Try to extract JSON if wrapped in text
if '{' in response_clean:
json_start = response_clean.find('{')
json_end = response_clean.rfind('}') + 1
response_clean = response_clean[json_start:json_end]
keywords_data = json.loads(response_clean)
# Ensure all required fields exist
result = {
"core_keywords": keywords_data.get("core_keywords", []),
"entities": keywords_data.get("entities", []),
"synonyms": keywords_data.get("synonyms", {}),
"expansions": keywords_data.get("expansions", {})
}
logger.info(f"Extracted keywords: {result['core_keywords'][:5]}, synonyms: {len(result['synonyms'])} terms")
return result
except json.JSONDecodeError as e:
logger.warning(f"Failed to parse LLM keyword extraction: {e}, using fallback")
# Fallback to simple extraction
words = query.split()
return {
"core_keywords": words,
"entities": [],
"synonyms": {},
"expansions": {}
}
except Exception as e:
logger.error(f"Keyword extraction failed: {e}", exc_info=True)
return {
"core_keywords": query.split(),
"entities": [],
"synonyms": {},
"expansions": {}
}
async def _retrieve_parallel(
self,
query: str,
user: str,
config: HybridRAGConfig,
keywords_data: Dict[str, Any]
) -> Dict[str, List]:
"""
Phase 1: Retrieve results from all sources in parallel.
Args:
query: Search query
user: User identifier
config: Search configuration
keywords_data: Extracted keywords/synonyms
Returns:
Dictionary with results from each source and timing
"""
tasks = {}
timing = {}
# Vector search
if config.enable_vector:
async def vector_search():
start = time.time()
try:
response = await self.vector.search(
query=query,
user=user,
limit=config.vector_limit
)
results = [
{
"page_id": r.page_id,
"title": r.page_title,
"content": r.content,
"path": r.page_path,
"score": r.score,
"source": "vector"
}
for r in response.results
]
return results, (time.time() - start) * 1000
except Exception as e:
logger.error(f"Vector search failed: {e}", exc_info=True)
return [], (time.time() - start) * 1000
tasks["vector"] = vector_search()
# Graph search
if config.enable_graph:
async def graph_search():
start = time.time()
try:
results = await self.graph.search_documents(
query=query,
user=user,
limit=config.graph_limit,
keywords_data=keywords_data
)
formatted = [
{
"page_id": r["page_id"],
"title": r["title"],
"content": "", # Graph doesn't return content
"path": r["path"],
"entity_matches": r.get("entity_matches", 0),
"matched_entities": r.get("matched_entities", []),
"source": "graph"
}
for r in results
]
return formatted, (time.time() - start) * 1000
except Exception as e:
logger.error(f"Graph search failed: {e}", exc_info=True)
return [], (time.time() - start) * 1000
tasks["graph"] = graph_search()
# Web search
if config.enable_web:
async def web_search():
start = time.time()
try:
results = await self.searxng.search_general(
query=query,
limit=config.web_limit
)
formatted = [
{
"url": r.get("url"),
"title": r.get("title", ""),
"content": r.get("content", ""),
"engine": r.get("engine", ""),
"source": "web"
}
for r in results
]
return formatted, (time.time() - start) * 1000
except Exception as e:
logger.error(f"Web search failed: {e}", exc_info=True)
return [], (time.time() - start) * 1000
tasks["web"] = web_search()
# Execute all searches in parallel
results_dict = await asyncio.gather(*tasks.values())
# Combine results with timing
output = {"timing": {}}
for i, source in enumerate(tasks.keys()):
results, source_timing = results_dict[i]
output[source] = results
output["timing"][f"{source}_ms"] = source_timing
logger.info(
f"Parallel retrieval: vector={len(output.get('vector', []))}, "
f"graph={len(output.get('graph', []))}, web={len(output.get('web', []))}"
)
return output
def _reciprocal_rank_fusion(
self,
results_by_source: Dict[str, List],
k: int = 60
) -> List[Dict[str, Any]]:
"""
Phase 2: Merge results using Reciprocal Rank Fusion.
RRF formula: score = sum(1 / (k + rank)) for each source
Args:
results_by_source: Results from each source
k: RRF constant (default 60)
Returns:
Merged and sorted results
"""
rrf_scores = {}
for source, results in results_by_source.items():
for rank, result in enumerate(results, start=1):
# Use page_id for wiki results, url hash for web results
if result.get("page_id"):
result_id = f"page_{result['page_id']}"
elif result.get("url"):
result_id = f"url_{hash(result['url'])}"
else:
continue # Skip results without ID
if result_id not in rrf_scores:
rrf_scores[result_id] = {
"result": result,
"rrf_score": 0.0,
"sources": [],
"source_type": source
}
# RRF formula: sum of 1/(k + rank) across sources
rrf_scores[result_id]["rrf_score"] += 1 / (k + rank)
rrf_scores[result_id]["sources"].append(source)
# If result appears in multiple sources, update source_type
if len(rrf_scores[result_id]["sources"]) > 1:
rrf_scores[result_id]["source_type"] = "+".join(
sorted(set(rrf_scores[result_id]["sources"]))
)
# Sort by RRF score descending
sorted_results = sorted(
rrf_scores.values(),
key=lambda x: x["rrf_score"],
reverse=True
)
logger.info(f"RRF fusion: {len(sorted_results)} unique results from {len(results_by_source)} sources")
return sorted_results
async def _enrich_with_related_dossiers(
self,
results: List[Dict[str, Any]],
user: str
) -> List[Dict[str, Any]]:
"""
Phase 3: Enrich results with related documents via shared entities.
Args:
results: Fused results
user: User identifier
Returns:
Results with related_dossiers added
"""
for result in results:
result_data = result.get("result", {})
page_id = result_data.get("page_id")
if page_id:
try:
related_docs = await self.graph.get_related_documents(
page_id=page_id,
user=user,
limit=5
)
# Convert to RelatedDossier format
related_dossiers = []
for doc in related_docs:
for tag in doc.get("tags", [])[:3]: # Max 3 tags per doc
related_dossiers.append({
"page_id": doc["page_id"],
"title": doc["title"],
"path": doc["path"],
"tag": tag,
"shared_entities": doc["shared_entities"]
})
result["related_dossiers"] = related_dossiers[:5] # Limit to 5 total
except Exception as e:
logger.warning(f"Failed to get related docs for page {page_id}: {e}")
result["related_dossiers"] = []
else:
result["related_dossiers"] = []
return results
async def _rerank_with_llm(
self,
results: List[Dict[str, Any]],
query: str
) -> List[Dict[str, Any]]:
"""
Phase 4: Re-rank results using LLM for better relevance.
Args:
results: Results to re-rank (top 20)
query: Original search query
Returns:
Re-ranked results
"""
if len(results) <= 1:
return results
try:
# Build prompt with numbered results
docs_text = "\n".join([
f"{i+1}. {r['result'].get('title', 'Untitled')} - {r['result'].get('content', '')[:200]}..."
for i, r in enumerate(results)
])
prompt = f"""Given this search query and documents, rank them by relevance.
Query: {query}
Documents:
{docs_text}
Return only the numbers in order of relevance (most relevant first).
Example: 3,1,5,2,4
Ranking:"""
response = await self.ollama.generate_text(
prompt=prompt,
model=self.reranker_model
)
# Parse response: "3,1,5,2,4" → [2, 0, 4, 1, 3] (0-indexed)
indices_str = response.strip().split('\n')[0] # Take first line
indices = [int(x.strip()) - 1 for x in indices_str.split(",") if x.strip().isdigit()]
# Reorder results according to LLM ranking
reranked = []
for idx in indices:
if 0 <= idx < len(results):
reranked.append(results[idx])
# Add any results that weren't in the LLM response
for i, result in enumerate(results):
if i not in indices and result not in reranked:
reranked.append(result)
logger.info(f"LLM re-ranking: reordered {len(reranked)} results")
return reranked
except Exception as e:
logger.warning(f"LLM re-ranking failed: {e}, using RRF order")
return results # Fallback to RRF order
def _format_context_for_llm(self, results: List[HybridRAGResult]) -> str:
"""
Phase 5: Format results into context for LLM consumption.
Args:
results: Ranked results
Returns:
Formatted context string
"""
context_parts = []
for i, result in enumerate(results[:10], start=1):
# Source indicator
source_tag = f"[{result.source_type.upper()}]"
# Related dossiers if available
related = ""
if result.related_dossiers:
tags = ", ".join([d.tag for d in result.related_dossiers[:3]])
related = f"\n Related research: {tags}"
# Build context entry
content_preview = result.content[:300] if result.content else "(no content)"
context_parts.append(
f"{i}. {source_tag} {result.title}\n"
f" {content_preview}...{related}"
)
return "\n\n".join(context_parts)
async def _persist_search_for_librarian(
self,
query: str,
user: str,
keywords_data: Dict[str, Any],
raw_results: Dict[str, List],
final_results: List[Dict[str, Any]],
timing: Dict[str, float]
) -> Optional[str]:
"""
Phase 6: Store search query and results for Librarian processing.
Creates SearchQuery node in Neo4j with relationships to found documents
and web results for offline knowledge consolidation.
Args:
query: Search query
user: User identifier
keywords_data: Extracted keywords/synonyms
raw_results: Results from each source
final_results: Final ranked results
timing: Performance timing
Returns:
Search ID for tracking
"""
try:
user_base_label = get_neo4j_user_base_label(user)
search_id = str(uuid.uuid4())
# Create SearchQuery node
create_query = f"""
CREATE (sq:{user_base_label}_SearchQuery:SearchQuery {{
id: $search_id,
query: $query,
user: $user,
timestamp: datetime(),
processed: false,
total_results: $total_results,
vector_count: $vector_count,
graph_count: $graph_count,
web_count: $web_count,
keywords: $keywords,
synonyms: $synonyms,
timing_ms: $timing_ms
}})
RETURN sq.id as id
"""
result = await self.graph.neo4j.execute_query(create_query, {
"search_id": search_id,
"query": query,
"user": user,
"total_results": len(final_results),
"vector_count": len(raw_results.get("vector", [])),
"graph_count": len(raw_results.get("graph", [])),
"web_count": len(raw_results.get("web", [])),
"keywords": keywords_data.get("core_keywords", []),
"synonyms": json.dumps(keywords_data.get("synonyms", {})),
"timing_ms": timing.get("total_ms", 0)
})
# Link to found wiki documents (top 20)
for rank, result_data in enumerate(final_results[:20], start=1):
result = result_data.get("result", {})
page_id = result.get("page_id")
if page_id:
link_doc_query = f"""
MATCH (sq:{user_base_label}_SearchQuery:SearchQuery {{id: $search_id}})
MATCH (d:Document {{page_id: $page_id}})
MERGE (sq)-[f:FOUND]->(d)
SET f.source = $source,
f.rank = $rank,
f.rrf_score = $rrf_score,
f.final_rank = $final_rank
"""
await self.graph.neo4j.execute_query(link_doc_query, {
"search_id": search_id,
"page_id": page_id,
"source": result_data.get("source_type", "unknown"),
"rank": rank,
"rrf_score": result_data.get("rrf_score", 0),
"final_rank": result_data.get("final_rank", rank)
})
# Store web results as WebResult nodes (top 10)
web_results = [r for r in final_results[:10] if r.get("result", {}).get("url")]
for rank, result_data in enumerate(web_results, start=1):
result = result_data.get("result", {})
create_web_query = f"""
MATCH (sq:{user_base_label}_SearchQuery:SearchQuery {{id: $search_id}})
CREATE (wr:{user_base_label}_WebResult:WebResult {{
url: $url,
title: $title,
content: $content,
search_id: $search_id,
timestamp: datetime()
}})
CREATE (sq)-[:FOUND {{
source: "web",
rank: $rank,
rrf_score: $rrf_score
}}]->(wr)
"""
await self.graph.neo4j.execute_query(create_web_query, {
"search_id": search_id,
"url": result.get("url"),
"title": result.get("title", ""),
"content": result.get("content", "")[:1000], # Truncate
"rank": rank,
"rrf_score": result_data.get("rrf_score", 0)
})
logger.info(f"Persisted search {search_id} for Librarian processing")
return search_id
except Exception as e:
logger.error(f"Failed to persist search for Librarian: {e}", exc_info=True)
return None
def _convert_to_result_models(self, results: List[Dict[str, Any]]) -> List[HybridRAGResult]:
"""
Convert internal result format to HybridRAGResult models.
Args:
results: Internal result dictionaries
Returns:
List of HybridRAGResult models
"""
models = []
for result_data in results:
result = result_data.get("result", {})
related_dossiers = result_data.get("related_dossiers", [])
models.append(HybridRAGResult(
source_type=result_data.get("source_type", "unknown"),
title=result.get("title", "Untitled"),
content=result.get("content", ""),
url=result.get("url"),
page_id=result.get("page_id"),
page_path=result.get("path"),
rrf_score=result_data.get("rrf_score", 0),
final_rank=result_data.get("final_rank", 0),
sources=result_data.get("sources", []),
related_dossiers=[RelatedDossier(**d) for d in related_dossiers],
metadata={
"entity_matches": result.get("entity_matches"),
"matched_entities": result.get("matched_entities"),
"engine": result.get("engine")
}
))
return models
+414
View File
@@ -0,0 +1,414 @@
"""
Document Ingestion Service
Orchestrates the ingestion of wiki pages into the knowledge base:
1. Fetches page content from Wiki.js
2. Generates vector embeddings (Qdrant)
3. Extracts entities and updates knowledge graph (Neo4j)
This service is called by:
- Consolidation service (after creating/updating pages)
- Manual ingestion endpoints
- Batch ingestion jobs
"""
import logging
import asyncio
from typing import List, Optional
from datetime import datetime
import time
from src.services.vector_service import VectorService
from src.services.graph_service import GraphService
from src.clients.wikijs_client import WikiJSClient
from src.models.ingestion import (
IngestionRequest,
IngestionResult,
BatchIngestionRequest,
BatchIngestionResult
)
logger = logging.getLogger(__name__)
class IngestionService:
"""
Service for ingesting wiki pages into the knowledge base.
"""
def __init__(
self,
vector_service: VectorService,
graph_service: GraphService,
wiki_client: WikiJSClient
):
self.vector = vector_service
self.graph = graph_service
self.wiki = wiki_client
async def ingest_page(
self,
page_id: int,
user: str,
force_refresh: bool = False,
skip_vectors: bool = False,
skip_graph: bool = False,
skip_entity_linking: bool = False
) -> IngestionResult:
"""
Ingest a single wiki page into the knowledge base.
Args:
page_id: Wiki page ID
user: User identifier
force_refresh: Force re-ingestion even if unchanged
skip_vectors: Skip vector embedding generation
skip_graph: Skip graph entity extraction
skip_entity_linking: Skip automatic entity linking
Returns:
IngestionResult with operation details
"""
start_time = time.time()
logger.info(f"Starting ingestion for page {page_id} (user: {user})")
try:
# Fetch page to get metadata
page = await self.wiki.get_page(page_id)
if not page:
return IngestionResult(
page_id=page_id,
page_title=f"Page {page_id}",
success=False,
error="Page not found in Wiki.js",
processing_time_ms=(time.time() - start_time) * 1000
)
page_title = page.get("title", f"Page {page_id}")
page_path = page.get("path", "")
# Ingest vectors and graph in parallel
tasks = []
if not skip_vectors:
tasks.append(self._ingest_vectors(page_id, user, force_refresh))
else:
tasks.append(asyncio.create_task(asyncio.sleep(0))) # Dummy task
if not skip_graph:
tasks.append(self._ingest_graph(page_id, user, force_refresh))
else:
tasks.append(asyncio.create_task(asyncio.sleep(0))) # Dummy task
# Execute in parallel
vector_result, graph_result = await asyncio.gather(*tasks, return_exceptions=True)
# Handle errors
vector_chunks = 0
graph_entities = 0
graph_relationships = 0
errors = []
if not skip_vectors:
if isinstance(vector_result, Exception):
errors.append(f"Vector ingestion failed: {str(vector_result)}")
logger.error(f"Vector ingestion failed for page {page_id}: {vector_result}")
else:
vector_chunks = vector_result.get("chunks_created", 0)
if not skip_graph:
if isinstance(graph_result, Exception):
errors.append(f"Graph ingestion failed: {str(graph_result)}")
logger.error(f"Graph ingestion failed for page {page_id}: {graph_result}")
else:
# entities_extracted is a list, get its length
entities_list = graph_result.get("entities_extracted", [])
graph_entities = len(entities_list) if isinstance(entities_list, list) else 0
graph_relationships = graph_result.get("relationships_created", 0)
# Step 3: Link existing entities in the page content (after graph extraction)
entity_links_created = 0
if not skip_entity_linking and not skip_graph and not isinstance(graph_result, Exception):
try:
entity_links_created = await self._link_existing_entities(page_id, user, page)
logger.info(f"Created {entity_links_created} entity mention links for page {page_id}")
except Exception as e:
logger.warning(f"Entity linking failed for page {page_id}: {e}")
# Don't fail the whole ingestion if entity linking fails
processing_time_ms = (time.time() - start_time) * 1000
result = IngestionResult(
page_id=page_id,
page_title=page_title,
page_path=page_path,
success=len(errors) == 0,
error="; ".join(errors) if errors else None,
vector_chunks_created=vector_chunks,
graph_entities_extracted=graph_entities,
graph_relationships_created=graph_relationships,
processing_time_ms=processing_time_ms
)
if result.success:
logger.info(
f"Successfully ingested page {page_id}: "
f"{vector_chunks} chunks, {graph_entities} entities, "
f"{graph_relationships} relationships, {entity_links_created} entity links "
f"in {processing_time_ms:.0f}ms"
)
else:
logger.warning(f"Partial ingestion failure for page {page_id}: {result.error}")
return result
except Exception as e:
logger.error(f"Ingestion failed for page {page_id}: {e}", exc_info=True)
return IngestionResult(
page_id=page_id,
page_title=f"Page {page_id}",
success=False,
error=str(e),
processing_time_ms=(time.time() - start_time) * 1000
)
async def _ingest_vectors(
self,
page_id: int,
user: str,
force_refresh: bool
) -> dict:
"""
Ingest page into vector database.
Returns:
Dict with chunks_created count
"""
try:
summary = await self.vector.update_from_page(
page_id=page_id,
user=user,
force_refresh=force_refresh
)
return {
"chunks_created": summary.chunks_created,
"chunks_deleted": summary.chunks_deleted
}
except Exception as e:
logger.error(f"Vector ingestion failed for page {page_id}: {e}")
raise
async def _ingest_graph(
self,
page_id: int,
user: str,
force_refresh: bool
) -> dict:
"""
Ingest page into knowledge graph.
Returns:
Dict with entities_extracted and relationships_created counts
"""
try:
summary = await self.graph.update_from_page(
page_id=page_id,
user=user,
force_refresh=force_refresh
)
return {
"entities_extracted": summary.entities_extracted,
"relationships_created": summary.relationships_created
}
except Exception as e:
logger.error(f"Graph ingestion failed for page {page_id}: {e}")
raise
async def _link_existing_entities(
self,
page_id: int,
user: str,
page: dict
) -> int:
"""
Find and link mentions of existing entities in the page content.
This runs automatically after graph extraction to create MENTIONS relationships
for entities that already exist in the knowledge graph but were mentioned in
this page.
Args:
page_id: Wiki page ID
user: User identifier
page: Page dict with content (from WikiJSClient)
Returns:
Number of new entity mention links created
"""
import re
try:
page_content = page.get("content", "")
if not page_content or len(page_content) < 10:
return 0
# Get all existing entities from the knowledge graph
entities = await self.graph.get_all_entities(user)
if not entities:
logger.debug(f"No existing entities found for user {user}, skipping entity linking")
return 0
# Find entity mentions in page content
found_entities = []
content_lower = page_content.lower()
for entity in entities:
entity_name = entity.get("name", "")
if not entity_name or len(entity_name) < 3:
continue
# Create regex pattern for whole word matching
# This avoids matching "John" in "Johnson"
pattern = r'\b' + re.escape(entity_name.lower()) + r'\b'
# Find all matches
matches = list(re.finditer(pattern, content_lower))
if matches:
found_entities.append({
"name": entity_name,
"type": entity.get("type", "unknown"),
"mentions": len(matches),
"entity_id": entity.get("id")
})
if not found_entities:
logger.debug(f"No entity mentions found in page {page_id}")
return 0
# Create MENTIONS relationships
new_links_created = await self.graph.create_entity_mentions(
page_id=page_id,
user=user,
entity_names=found_entities
)
return new_links_created
except Exception as e:
logger.error(f"Entity linking failed for page {page_id}: {e}")
raise
async def ingest_batch(
self,
page_ids: List[int],
user: str,
force_refresh: bool = False,
skip_vectors: bool = False,
skip_graph: bool = False,
max_concurrent: int = 3
) -> BatchIngestionResult:
"""
Ingest multiple wiki pages concurrently.
Args:
page_ids: List of wiki page IDs to ingest
user: User identifier
force_refresh: Force re-ingestion
skip_vectors: Skip vector embedding generation
skip_graph: Skip graph entity extraction
max_concurrent: Maximum concurrent ingestion tasks
Returns:
BatchIngestionResult with per-page results
"""
start_time = time.time()
logger.info(f"Starting batch ingestion of {len(page_ids)} pages (user: {user})")
results = []
semaphore = asyncio.Semaphore(max_concurrent)
async def ingest_with_semaphore(page_id: int):
async with semaphore:
return await self.ingest_page(
page_id=page_id,
user=user,
force_refresh=force_refresh,
skip_vectors=skip_vectors,
skip_graph=skip_graph
)
# Execute all ingestions with concurrency control
tasks = [ingest_with_semaphore(page_id) for page_id in page_ids]
results = await asyncio.gather(*tasks)
# Calculate summary
successful = sum(1 for r in results if r.success)
failed = len(results) - successful
total_processing_time_ms = (time.time() - start_time) * 1000
batch_result = BatchIngestionResult(
total_pages=len(page_ids),
successful=successful,
failed=failed,
results=results,
total_processing_time_ms=total_processing_time_ms
)
logger.info(
f"Batch ingestion complete: {successful}/{len(page_ids)} successful "
f"in {total_processing_time_ms:.0f}ms"
)
return batch_result
async def ingest_all_pages(
self,
user: str,
path_prefix: Optional[str] = None,
force_refresh: bool = False,
max_concurrent: int = 3
) -> BatchIngestionResult:
"""
Ingest all wiki pages for a user.
Args:
user: User identifier
path_prefix: Optional path prefix filter (e.g., "users/jpmschweitzer")
force_refresh: Force re-ingestion
max_concurrent: Maximum concurrent ingestion tasks
Returns:
BatchIngestionResult
"""
logger.info(f"Finding all pages for user {user} (prefix: {path_prefix or 'all'})")
# List all pages (not search - search requires a query and may have stale index)
pages = await self.wiki.list_all_pages(
path_prefix=path_prefix or f"users/{user}"
)
if not pages:
logger.warning(f"No pages found for user {user}")
return BatchIngestionResult(
total_pages=0,
successful=0,
failed=0,
results=[],
total_processing_time_ms=0
)
page_ids = [p['id'] for p in pages]
logger.info(f"Found {len(page_ids)} pages to ingest")
return await self.ingest_batch(
page_ids=page_ids,
user=user,
force_refresh=force_refresh,
max_concurrent=max_concurrent
)
+358
View File
@@ -0,0 +1,358 @@
"""
Vector service for Library Desk Qdrant operations.
Handles semantic search, document chunking, and embeddings.
"""
import re
import time
import hashlib
import uuid
from typing import List, Dict, Any, Optional
import logging
from src.clients.qdrant_client import QdrantClientWrapper
from src.clients.wikijs_client import WikiJSClient
from src.clients.ollama_client import OllamaClient
from src.core.multi_tenancy import get_qdrant_collection_name
from src.models.vector import (
SearchResult, SearchResponse, VectorUpdateSummary,
DocumentChunk, CollectionInfo, CollectionListResponse
)
logger = logging.getLogger(__name__)
class VectorService:
"""
Service for Qdrant vector operations.
Responsibilities:
- Document chunking
- Embedding generation
- Semantic search
- Vector CRUD operations
"""
def __init__(
self,
qdrant_client: QdrantClientWrapper,
wikijs_client: WikiJSClient,
ollama_client: OllamaClient,
chunk_size: int = 500,
chunk_overlap: int = 50
):
"""
Initialize vector service.
Args:
qdrant_client: Qdrant database client
wikijs_client: Wiki.js client for fetching pages
ollama_client: Ollama client for embeddings
chunk_size: Target chunk size in tokens (approximate)
chunk_overlap: Overlap between chunks in tokens
"""
self.qdrant = qdrant_client
self.wiki = wikijs_client
self.ollama = ollama_client
self.chunk_size = chunk_size
self.chunk_overlap = chunk_overlap
def _chunk_text(self, text: str) -> List[str]:
"""
Chunk text into overlapping segments.
Simple word-based chunking for now.
TODO: Use tiktoken or similar for token-accurate chunking.
Args:
text: Text to chunk
Returns:
List of text chunks
"""
# Remove extra whitespace
text = re.sub(r'\s+', ' ', text).strip()
# Split into words (approximates tokens)
words = text.split()
if len(words) <= self.chunk_size:
return [text]
chunks = []
start = 0
while start < len(words):
end = start + self.chunk_size
chunk_words = words[start:end]
chunks.append(' '.join(chunk_words))
# Move start forward with overlap
start = end - self.chunk_overlap
return chunks
async def update_from_page(
self,
page_id: int,
user: str,
force_refresh: bool = False
) -> VectorUpdateSummary:
"""
Update vector embeddings from a wiki page.
Chunks the page content, generates embeddings, and upserts to Qdrant.
Args:
page_id: Wiki page ID
user: User identifier
force_refresh: Force re-embedding even if unchanged
Returns:
Summary of update operation
"""
start_time = time.time()
try:
# Fetch page from Wiki.js
page = await self.wiki.get_page(page_id)
if not page:
raise ValueError(f"Page {page_id} not found")
# Get collection name for user
collection_name = get_qdrant_collection_name(user)
# Ensure collection exists
await self.qdrant.ensure_collection(collection_name)
# Extract content
content = page.get("content", "")
title = page.get("title", "")
path = page.get("path", "")
if not content:
logger.warning(f"Page {page_id} has no content, skipping vector update")
return VectorUpdateSummary(
page_id=page_id,
page_title=title,
processing_time_ms=(time.time() - start_time) * 1000,
success=True
)
# Chunk the content
chunks = self._chunk_text(content)
logger.info(f"Split page {page_id} into {len(chunks)} chunks")
# Delete existing chunks for this page
deleted_count = await self.qdrant.delete_by_filter(
collection_name=collection_name,
filter_conditions={"page_id": page_id}
)
# Generate embeddings and upsert chunks
chunks_created = 0
for idx, chunk_text in enumerate(chunks):
# Generate deterministic UUID from page_id and chunk_index
chunk_id = str(uuid.uuid5(uuid.NAMESPACE_DNS, f"page_{page_id}_chunk_{idx}"))
# Generate embedding
embedding = await self.ollama.embed(chunk_text)
if not embedding:
logger.error(f"Failed to generate embedding for chunk {chunk_id}")
continue
# Prepare metadata
metadata = {
"page_id": page_id,
"page_title": title,
"page_path": path,
"chunk_index": idx,
"chunk_text": chunk_text,
"user": user
}
# Upsert to Qdrant
success = await self.qdrant.upsert_vector(
collection_name=collection_name,
vector_id=chunk_id,
vector=embedding,
payload=metadata
)
if success:
chunks_created += 1
processing_time_ms = (time.time() - start_time) * 1000
logger.info(
f"Updated vectors for page {page_id}: "
f"{chunks_created} chunks created, {deleted_count} old chunks deleted"
)
return VectorUpdateSummary(
page_id=page_id,
page_title=title,
chunks_created=chunks_created,
chunks_deleted=deleted_count,
total_chunks=chunks_created,
embedding_dim=len(embedding) if embedding else 768,
processing_time_ms=processing_time_ms,
success=True
)
except Exception as e:
processing_time_ms = (time.time() - start_time) * 1000
logger.error(f"Failed to update vectors for page {page_id}: {e}", exc_info=True)
return VectorUpdateSummary(
page_id=page_id,
page_title="Unknown",
processing_time_ms=processing_time_ms,
success=False,
error_message=str(e)
)
async def search(
self,
query: str,
user: str,
limit: int = 10,
score_threshold: float = 0.5
) -> SearchResponse:
"""
Perform semantic search across user's documents.
Args:
query: Search query text
user: User identifier
limit: Maximum results to return
score_threshold: Minimum similarity score (0-1)
Returns:
Search results with similarity scores
"""
start_time = time.time()
try:
# Get collection name
collection_name = get_qdrant_collection_name(user)
# Check if collection exists
exists = await self.qdrant.collection_exists(collection_name)
if not exists:
logger.info(f"Collection {collection_name} doesn't exist, returning empty results")
return SearchResponse(
query=query,
results=[],
total=0,
user=user
)
# Generate query embedding
query_embedding = await self.ollama.embed(query)
if not query_embedding:
raise ValueError("Failed to generate query embedding")
# Search in Qdrant
search_results = await self.qdrant.search_vectors(
collection_name=collection_name,
query_vector=query_embedding,
limit=limit,
score_threshold=score_threshold
)
# Convert to SearchResult models
results = []
for result in search_results:
payload = result.get("payload", {})
results.append(SearchResult(
chunk_id=result["id"],
page_id=payload.get("page_id", 0),
page_title=payload.get("page_title"),
page_path=payload.get("page_path"),
chunk_index=payload.get("chunk_index", 0),
content=payload.get("chunk_text", ""),
score=result["score"],
metadata=payload
))
query_time_ms = (time.time() - start_time) * 1000
logger.info(f"Semantic search completed in {query_time_ms:.2f}ms: {len(results)} results")
return SearchResponse(
query=query,
results=results,
total=len(results),
user=user
)
except Exception as e:
logger.error(f"Semantic search failed: {e}", exc_info=True)
return SearchResponse(
query=query,
results=[],
total=0,
user=user
)
async def delete_page_chunks(
self,
page_id: int,
user: str
) -> int:
"""
Delete all chunks for a wiki page.
Args:
page_id: Wiki page ID
user: User identifier
Returns:
Number of chunks deleted
"""
collection_name = get_qdrant_collection_name(user)
try:
deleted_count = await self.qdrant.delete_by_filter(
collection_name=collection_name,
filter_conditions={"page_id": page_id}
)
logger.info(f"Deleted {deleted_count} chunks for page {page_id}")
return deleted_count
except Exception as e:
logger.error(f"Failed to delete chunks for page {page_id}: {e}", exc_info=True)
return 0
async def list_collections(self) -> CollectionListResponse:
"""
List all Qdrant collections.
Returns:
List of collections with stats
"""
try:
collections_data = await self.qdrant.list_collections()
collections = []
for coll in collections_data:
collections.append(CollectionInfo(
name=coll["name"],
vectors_count=coll.get("vectors_count", 0),
points_count=coll.get("points_count", 0),
segments_count=coll.get("segments_count", 0)
))
return CollectionListResponse(
collections=collections,
total=len(collections)
)
except Exception as e:
logger.error(f"Failed to list collections: {e}", exc_info=True)
return CollectionListResponse(
collections=[],
total=0
)
+279
View File
@@ -0,0 +1,279 @@
"""
Wiki.js Database Change Listener
Listens to PostgreSQL NOTIFY events for page changes in Wiki.js
and triggers the same processing as webhooks would.
This is an alternative to Wiki.js webhooks (which don't exist in open-source version).
"""
import logging
import asyncio
import asyncpg
from typing import Optional
from datetime import datetime
from src.config import get_settings
from src.core.dependencies import get_ingestion_service
from src.services.consolidation_service import ConsolidationService
logger = logging.getLogger(__name__)
class WikiChangeListener:
"""
Listens to PostgreSQL NOTIFY events from Wiki.js database.
This requires setting up triggers in the Wiki.js database to emit
NOTIFY events on INSERT/UPDATE/DELETE to the pages table.
"""
def __init__(self):
self.settings = get_settings()
self.connection: Optional[asyncpg.Connection] = None
self.running = False
# Loop prevention: Track recently processed pages
# Key: page_id, Value: timestamp of last processing
self._recent_notifications = {}
self._debounce_seconds = self.settings.wikijs_change_listener_debounce_seconds
async def start(self):
"""Start listening to database changes."""
logger.info("Starting Wiki.js database change listener")
# Connect to Wiki.js PostgreSQL database
self.connection = await asyncpg.connect(
host=self.settings.wikijs_db_host,
port=self.settings.wikijs_db_port,
user=self.settings.wikijs_db_user,
password=self.settings.wikijs_db_password,
database=self.settings.wikijs_db_name
)
# Listen to the wiki_page_changes channel
await self.connection.add_listener('wiki_page_changes', self._handle_notification)
self.running = True
logger.info("Listening for Wiki.js page changes via PostgreSQL NOTIFY")
async def stop(self):
"""Stop listening and close connection."""
if self.connection:
await self.connection.remove_listener('wiki_page_changes', self._handle_notification)
await self.connection.close()
self.running = False
logger.info("Stopped Wiki.js change listener")
async def _handle_notification(self, connection, pid, channel, payload):
"""Handle NOTIFY event from database."""
try:
# Payload format: "operation:page_id:user_email"
# e.g., "INSERT:123:user@example.com"
parts = payload.split(':')
if len(parts) < 3:
logger.warning(f"Invalid notification payload: {payload}")
return
operation = parts[0] # INSERT, UPDATE, DELETE
page_id = int(parts[1])
user_email = parts[2]
logger.info(f"Received {operation} notification for page {page_id} by {user_email}")
# LOOP PREVENTION: Debouncing - ignore rapid duplicate notifications
# Note: We rely solely on debouncing for loop prevention because:
# - The user_email in notifications is the page creator, not the editor
# - Creator != namespace owner (e.g., 'librarian' creates page in 'users/jpmschweitzer/')
# - Filtering by creator breaks legitimate page ingestion
if self._is_recently_processed(page_id):
logger.debug(
f"Skipping notification for page {page_id} - "
f"processed within last {self._debounce_seconds}s (debouncing)"
)
return
# Mark as recently processed
self._mark_as_processed(page_id)
# Map operation to webhook-style event
event_map = {
'INSERT': 'page.create',
'UPDATE': 'page.update',
'DELETE': 'page.delete'
}
event = event_map.get(operation, 'page.update')
# Extract user from email
user = user_email.split('@')[0] if '@' in user_email else 'jpmschweitzer'
# Process the change
await self._process_page_change(
page_id=page_id,
event=event,
user=user
)
except Exception as e:
logger.error(f"Failed to handle notification: {e}", exc_info=True)
def _is_automated_user(self, email: str) -> bool:
"""
Check if email belongs to an automated system user.
These are edits made by library-desk via Wiki.js API (entity linking).
We skip processing these to prevent loops.
Customize this list based on your Wiki.js username for library-desk.
"""
automated_users = [
self.settings.wikijs_username, # Library-desk's Wiki.js API user
"library-desk@system",
"automation@system",
"bot@system"
]
return email.lower() in [u.lower() for u in automated_users]
def _is_recently_processed(self, page_id: int) -> bool:
"""Check if page was processed recently (debouncing)."""
if page_id not in self._recent_notifications:
return False
last_processed = self._recent_notifications[page_id]
elapsed = (datetime.now() - last_processed).total_seconds()
return elapsed < self._debounce_seconds
def _mark_as_processed(self, page_id: int):
"""Mark page as recently processed."""
self._recent_notifications[page_id] = datetime.now()
# Clean up old entries (keep last 100 pages)
if len(self._recent_notifications) > 100:
# Remove oldest entries
sorted_items = sorted(
self._recent_notifications.items(),
key=lambda x: x[1]
)
self._recent_notifications = dict(sorted_items[-100:])
async def _process_page_change(self, page_id: int, event: str, user: str):
"""Process page change identically to webhook handler."""
from src.routers.webhooks import process_wiki_page_change, cleanup_deleted_page
ingestion_service = get_ingestion_service()
if event == 'page.delete':
# For deletions, need to handle cleanup
# Note: We don't have page_title at this point, use page_id
await cleanup_deleted_page(
page_id=page_id,
page_title=f"Page {page_id}",
user=user,
ingestion_service=ingestion_service
)
else:
# For create/update, get page details and process
from src.core.dependencies import get_wiki_service
wiki_service = get_wiki_service()
try:
page = await wiki_service.get_page(page_id, user)
# If page access failed (wrong user), try to extract correct user from page path
if not page:
# Try to get page metadata without user validation to find correct namespace
try:
# Query Wiki.js directly for page path
page_info = await wiki_service.wiki_client.get_page(page_id)
if page_info and page_info.get('path'):
# Extract user from path: users/{user}/...
path_parts = page_info['path'].split('/')
if len(path_parts) >= 2 and path_parts[0] == 'users':
correct_user = path_parts[1]
logger.debug(f"Retrying page {page_id} with correct user: {correct_user}")
page = await wiki_service.get_page(page_id, correct_user)
user = correct_user
except Exception as e:
logger.debug(f"Could not extract user from page {page_id} path: {e}")
if page:
await process_wiki_page_change(
page_id=page_id,
page_title=page.title,
user=user,
event=event,
ingestion_service=ingestion_service
)
else:
logger.warning(f"Could not retrieve page {page_id} for processing")
except Exception as e:
logger.error(f"Failed to process page {page_id}: {e}")
# SQL to set up triggers in Wiki.js database
SETUP_TRIGGERS_SQL = """
-- Create function to notify on page changes
-- Note: Wiki.js pages table has authorId (FK to users.id), not authorEmail
-- We look up the email from the users table
CREATE OR REPLACE FUNCTION notify_page_change()
RETURNS TRIGGER AS $$
DECLARE
author_email TEXT;
BEGIN
IF TG_OP = 'DELETE' THEN
-- Look up email from users table using OLD.authorId
SELECT email INTO author_email FROM users WHERE id = OLD."authorId";
IF author_email IS NULL THEN
author_email := 'unknown@system';
END IF;
PERFORM pg_notify(
'wiki_page_changes',
TG_OP || ':' || OLD.id || ':' || author_email
);
RETURN OLD;
ELSE
-- Look up email from users table using NEW.authorId
SELECT email INTO author_email FROM users WHERE id = NEW."authorId";
IF author_email IS NULL THEN
author_email := 'unknown@system';
END IF;
PERFORM pg_notify(
'wiki_page_changes',
TG_OP || ':' || NEW.id || ':' || author_email
);
RETURN NEW;
END IF;
END;
$$ LANGUAGE plpgsql;
-- Create triggers on pages table
DROP TRIGGER IF EXISTS wiki_page_insert_trigger ON pages;
CREATE TRIGGER wiki_page_insert_trigger
AFTER INSERT ON pages
FOR EACH ROW
EXECUTE FUNCTION notify_page_change();
DROP TRIGGER IF EXISTS wiki_page_update_trigger ON pages;
CREATE TRIGGER wiki_page_update_trigger
AFTER UPDATE ON pages
FOR EACH ROW
EXECUTE FUNCTION notify_page_change();
DROP TRIGGER IF EXISTS wiki_page_delete_trigger ON pages;
CREATE TRIGGER wiki_page_delete_trigger
AFTER DELETE ON pages
FOR EACH ROW
EXECUTE FUNCTION notify_page_change();
-- Verify triggers are created
SELECT
trigger_name,
event_manipulation,
event_object_table
FROM information_schema.triggers
WHERE event_object_table = 'pages'
ORDER BY trigger_name;
"""
+493
View File
@@ -0,0 +1,493 @@
"""
Intelligent Wiki Page Writer Service
Uses LLM (mistral-nemo) to create and reconstruct wiki pages with:
- Holistic content restructuring
- Zero fact loss (unless superseded)
- Conflict detection and flagging
- Standard formatting with template adherence
- Professional organization (summary, tables, chapters)
This service is used by:
- Consolidation service (Librarian knowledge consolidation)
- Any other service that needs to create/update wiki pages
"""
import logging
import json
from typing import Dict, Any, List, Optional, Tuple
from datetime import datetime
logger = logging.getLogger(__name__)
class WikiPageWriter:
"""
Intelligent wiki page writer using LLM for content generation and restructuring.
"""
def __init__(self, ollama_client):
"""
Initialize wiki page writer.
Args:
ollama_client: OllamaClient for LLM operations
"""
self.ollama = ollama_client
self.model = "mistral-nemo" # Default model for writing
async def create_page(
self,
title: str,
topic_summary: str,
source_information: List[Dict[str, str]],
entities: Optional[List[str]] = None,
related_docs: Optional[List[str]] = None
) -> str:
"""
Create new wiki page with structured content.
Args:
title: Page title
topic_summary: Brief summary of the topic
source_information: List of {title, url, content} dicts
entities: Related entities from knowledge graph
related_docs: Related documents/pages
Returns:
Formatted markdown content
"""
logger.info(f"Creating wiki page: {title}")
# Build source context
sources_text = self._format_sources_for_llm(source_information)
# Create page using LLM
prompt = self._build_create_prompt(
title=title,
summary=topic_summary,
sources=sources_text,
entities=entities or [],
related_docs=related_docs or []
)
content = await self._call_llm(prompt)
# Post-process to ensure template compliance
content = self._ensure_standard_sections(
content=content,
title=title,
sources=source_information,
entities=entities or [],
related_docs=related_docs or []
)
return content
async def reconstruct_page(
self,
title: str,
existing_content: str,
new_information: str,
new_sources: List[Dict[str, str]],
detect_conflicts: bool = True
) -> Tuple[str, Optional[List[Dict[str, Any]]]]:
"""
Reconstruct existing page with new information.
Intelligently merges new content with existing, restructures for clarity,
and detects factual conflicts.
Args:
title: Page title
existing_content: Current page content
new_information: New information to integrate
new_sources: Sources for new information
detect_conflicts: Whether to detect and flag conflicts
Returns:
Tuple of (reconstructed_content, conflicts)
conflicts: List of detected conflicts or None
"""
logger.info(f"Reconstructing wiki page: {title}")
# Detect conflicts first
conflicts = None
if detect_conflicts:
conflicts = await self._detect_conflicts(
existing_content=existing_content,
new_information=new_information
)
if conflicts:
logger.warning(f"Detected {len(conflicts)} potential conflicts in {title}")
# Build reconstruction prompt
prompt = self._build_reconstruct_prompt(
title=title,
existing_content=existing_content,
new_information=new_information,
new_sources=self._format_sources_for_llm(new_sources),
conflicts=conflicts
)
# Reconstruct with LLM
reconstructed = await self._call_llm(prompt)
# Ensure standard sections are present
reconstructed = self._ensure_standard_sections(
content=reconstructed,
title=title,
sources=new_sources,
is_update=True
)
return reconstructed, conflicts
async def _detect_conflicts(
self,
existing_content: str,
new_information: str
) -> Optional[List[Dict[str, Any]]]:
"""
Detect factual conflicts between existing and new content.
Returns:
List of conflicts with: {fact_a, fact_b, confidence, context}
"""
prompt = f"""Analyze these two pieces of content for factual conflicts.
EXISTING CONTENT:
{existing_content[:2000]}
NEW INFORMATION:
{new_information[:2000]}
Identify any facts that contradict each other. For each conflict, provide:
1. The fact from existing content
2. The contradicting fact from new information
3. Confidence level (low/medium/high)
4. Context/explanation
Return ONLY valid JSON:
{{
"conflicts": [
{{
"existing_fact": "fact from old content",
"new_fact": "contradicting fact",
"confidence": "medium",
"context": "explanation of why these conflict"
}}
]
}}
If no conflicts, return: {{"conflicts": []}}
JSON:"""
try:
response = await self.ollama.generate_text(
prompt=prompt,
model=self.model,
stream=False
)
# Extract JSON
response_clean = response.strip()
if '{' in response_clean:
json_start = response_clean.find('{')
json_end = response_clean.rfind('}') + 1
response_clean = response_clean[json_start:json_end]
result = json.loads(response_clean)
conflicts = result.get('conflicts', [])
return conflicts if conflicts else None
except Exception as e:
logger.error(f"Conflict detection failed: {e}")
return None
def _build_create_prompt(
self,
title: str,
summary: str,
sources: str,
entities: List[str],
related_docs: List[str]
) -> str:
"""Build LLM prompt for creating new page."""
return f"""You are a Librarian creating a dossier for a personal knowledge base and extended memory system.
Create a comprehensive, well-structured wiki page with appropriate sections for the content type.
TOPIC: {title}
SUMMARY: {summary}
SOURCE INFORMATION:
{sources}
RELATED ENTITIES: {', '.join(entities) if entities else 'None'}
RELATED DOCUMENTS: {', '.join(related_docs) if related_docs else 'None'}
CONTENT TYPE GUIDELINES (Schema.org-aligned):
For PEOPLE (family, friends, colleagues, public figures) (Schema.org: Person):
- Executive Summary (who they are, key facts)
- Background & Biography
- Relationships & Connections
- Professional Info / Career
- Interests & Preferences
- Important Dates & Events
- Notes & Observations
For COMPANIES (businesses, organizations, startups) (Schema.org: Organization):
- Executive Summary (what they do, industry, key facts)
- Overview & Mission
- Products & Services
- History & Milestones
- Leadership & Team
- Personal Connection / Experience
- Notable Projects or Achievements
For PLACES (locations, restaurants, destinations) (Schema.org: Place):
- Executive Summary (what/where, key details)
- Location & How to Get There
- Description & Atmosphere
- Features & Amenities
- Personal Experiences / Visits
- Recommendations & Tips
For ENTERTAINMENT (books, movies, TV, music, games) (Schema.org: CreativeWork):
- Executive Summary (title, creator, key facts)
- Synopsis / Overview
- Key Characters / Themes
- Personal Thoughts & Ratings
- Memorable Moments / Quotes
- Related Works
For RECIPES & FOOD (Schema.org: Recipe):
- Executive Summary (dish name, cuisine type)
- Ingredients (formatted as table or list)
- Instructions (step-by-step)
- Cooking Tips & Variations
- Personal Notes & Modifications
- Source / Origin
For PRODUCTS (gear, tools, purchases) (Schema.org: Product):
- Executive Summary (what it is, brand/model, key specs)
- Overview & Purpose
- Specifications (formatted as table)
- Purchase Information (where, when, price)
- Personal Experience / Review
- Maintenance & Care
- Related Products / Alternatives
For TECHNOLOGY (software, applications, infrastructure) (Schema.org: SoftwareApplication):
- Executive Summary (what it is, key facts)
- Overview & Purpose
- Technical Details (tables for specs)
- Setup & Configuration
- Use Cases & Applications
- Best Practices
- Common Issues & Solutions
For EVENTS (concerts, travel, appointments) (Schema.org: Event):
- Executive Summary (what, when, where)
- Event Details (date, time, location, venue)
- Participants / Attendees
- Planning & Preparation
- Experience / Highlights
- Photos / Media
- Notes & Reflections
For HEALTH (medical, fitness, wellness) (Schema.org: MedicalEntity):
- Executive Summary (condition/topic, key facts)
- Overview & Background
- Symptoms / Signs / Characteristics
- Treatments / Approaches / Recommendations
- Personal Experience / Progress
- Resources & References
- Important Dates (appointments, changes)
For HOBBIES (activities, interests, pastimes) (Custom extension):
- Executive Summary (what it is, why interesting)
- Getting Started / Basics
- Equipment & Materials
- Techniques & Skills
- Personal Progress / Achievements
- Resources & Communities
- Goals & Future Plans
For PROJECTS (work projects, personal projects) (Schema.org: Project):
- Executive Summary (what, why, status)
- Goals & Objectives
- Timeline & Milestones
- Team / Collaborators
- Technical Details / Architecture
- Current Status & Next Steps
- Lessons Learned / Reflections
For REFERENCE (general knowledge, how-tos) (Custom extension):
- Executive Summary
- Overview & Context
- Key Concepts & Definitions
- Step-by-Step Guide (if applicable)
- Examples & Use Cases
- Tips & Best Practices
- Related Topics & Further Reading
FORMATTING RULES:
- Use markdown headers (##, ###)
- Create tables for structured data (ingredients, specs, comparisons)
- Use bullet points for lists
- Include code blocks with ``` where applicable
- Bold important terms
- Keep sections focused and scannable
- Adapt structure to content - not all sections apply to all topics
Generate ONLY the markdown content (do not include Sources, Knowledge Graph, or Mind Map sections - those are added automatically).
MARKDOWN:"""
def _build_reconstruct_prompt(
self,
title: str,
existing_content: str,
new_information: str,
new_sources: str,
conflicts: Optional[List[Dict[str, Any]]]
) -> str:
"""Build LLM prompt for reconstructing page."""
conflicts_note = ""
if conflicts:
conflicts_note = "\n\nDETECTED CONFLICTS:\n"
for i, c in enumerate(conflicts, 1):
conflicts_note += f"{i}. Existing: '{c['existing_fact']}'\n"
conflicts_note += f" New: '{c['new_fact']}'\n"
conflicts_note += f" Confidence: {c['confidence']}\n"
conflicts_note += f" Note: {c['context']}\n\n"
conflicts_note += "IMPORTANT: For conflicts, prefer the most recent/authoritative source. Add a note in 'Changes & Updates' section when facts are superseded.\n"
return f"""Reconstruct this wiki page by intelligently merging new information with existing content.
TITLE: {title}
EXISTING CONTENT:
{existing_content}
NEW INFORMATION TO INTEGRATE:
{new_information}
NEW SOURCES:
{new_sources}
{conflicts_note}
RECONSTRUCTION REQUIREMENTS:
1. **Zero Fact Loss**: Preserve ALL facts from existing content unless superseded
2. **Holistic Restructuring**: Reorganize for better flow and clarity
3. **Conflict Resolution**: When facts conflict, choose most authoritative/recent
4. **Professional Structure**:
- Update Executive Summary with key facts
- Organize into clear chapters
- Use tables for specifications/comparisons
- Maintain consistent formatting
5. **Update Tracking**: Add entry to "Changes & Updates" section with today's date
FORMATTING RULES:
- Maintain markdown structure
- Use tables for data (| col1 | col2 |)
- Keep existing good structure, improve where needed
- Bold important terms
- Add subsections (###) where it improves clarity
OUTPUT INSTRUCTIONS:
- Return complete page content (do not include Sources, Knowledge Graph, Mind Map - those are added automatically)
- Include updated "Changes & Updates" section noting what was changed today
- If facts were superseded, note it clearly
RECONSTRUCTED MARKDOWN:"""
async def _call_llm(self, prompt: str) -> str:
"""Call LLM with prompt and return response."""
try:
response = await self.ollama.generate_text(
prompt=prompt,
model=self.model,
stream=False
)
if not response:
raise Exception("Empty response from LLM")
return response.strip()
except Exception as e:
logger.error(f"LLM call failed: {e}")
raise
def _format_sources_for_llm(self, sources: List[Dict[str, str]]) -> str:
"""Format source information for LLM prompt."""
formatted = []
for i, source in enumerate(sources, 1):
formatted.append(f"[{i}] {source.get('title', 'Untitled')}")
formatted.append(f" URL: {source.get('url', 'N/A')}")
content = source.get('content', '')[:500] # Limit content length
formatted.append(f" Content: {content}...\n")
return "\n".join(formatted)
def _ensure_standard_sections(
self,
content: str,
title: str,
sources: List[Dict[str, str]],
entities: Optional[List[str]] = None,
related_docs: Optional[List[str]] = None,
is_update: bool = False
) -> str:
"""
Ensure page has standard footer sections (Sources, Knowledge Graph, Mind Map).
These sections are standardized and appended automatically.
"""
# Remove any existing standard sections
for section in ["## Sources", "## Knowledge Graph", "## Mind Map"]:
if section in content:
content = content.split(section)[0]
# Add horizontal rule before footer
content = content.rstrip() + "\n\n---\n\n"
# Add Sources section
content += "## Sources\n\n"
if sources:
for i, source in enumerate(sources, 1):
content += f"{i}. [{source.get('title', 'Source')}]({source.get('url', '#')})\n"
else:
content += "*No sources listed*\n"
# Add Knowledge Graph section
content += "\n## Knowledge Graph\n\n"
if entities:
content += "**Related Entities:**\n"
for entity in entities[:10]: # Limit to 10
content += f"- {entity}\n"
else:
content += "*No entities linked yet*\n"
content += "\n**View in Neo4j:** [Explore Graph](/graph)\n"
# Add Mind Map section
content += "\n## Mind Map\n\n"
content += f"**Interactive Mind Map:** [View Topic Map](/mindmap?topic={title.replace(' ', '+')})\n"
# Add footer metadata
content += "\n---\n\n"
timestamp = datetime.now().strftime('%Y-%m-%d %H:%M')
action = "Updated" if is_update else "Created"
content += f"*{action}: {timestamp} | Generated by: Librarian Agent* \n"
content += "*Template: Library Desk Wiki Standard v1.0*\n"
return content
+433
View File
@@ -0,0 +1,433 @@
"""
Wiki service layer for Library Desk.
Handles business logic for wiki operations with:
- Multi-tenant path scoping
- Dossier management (tag-based)
- Page CRUD operations
- Search functionality
"""
from typing import List, Optional, Dict, Any
import logging
from src.clients.wikijs_client import WikiJSClient
from src.core.multi_tenancy import get_wikijs_namespace, validate_user_id, DEFAULT_USER
from src.models.wiki import (
WikiPage, WikiPageSummary, WikiPageList,
WikiPageCreate, WikiPageUpdate,
DossierInfo, DossierList
)
logger = logging.getLogger(__name__)
class WikiService:
"""
Service layer for wiki operations.
Responsibilities:
- Enforce multi-tenant path scoping
- Convert between client and API models
- Handle dossier (tag) operations
- Provide business logic layer
"""
def __init__(self, wiki_client: WikiJSClient):
"""
Initialize wiki service.
Args:
wiki_client: Initialized Wiki.js client
"""
self.wiki_client = wiki_client
def _get_user_namespace(self, user: str) -> str:
"""
Get user's wiki namespace with validation.
Args:
user: User identifier
Returns:
Wiki.js namespace path
Raises:
ValueError: If user ID is invalid
"""
if not validate_user_id(user):
raise ValueError(f"Invalid user ID: {user}")
return get_wikijs_namespace(user)
def _ensure_user_path(self, path: str, user: str) -> str:
"""
Ensure path is within user's namespace.
Args:
path: Requested page path
user: User identifier
Returns:
Full path within user namespace
Example:
>>> self._ensure_user_path("/projects/foo", "jpmschweitzer")
'/users/jpmschweitzer/projects/foo'
"""
namespace = self._get_user_namespace(user)
# If path already starts with namespace, return as-is
if path.startswith(namespace):
return path
# Remove leading slash from path if present
path = path.lstrip("/")
# Combine namespace and path
return f"{namespace}/{path}"
async def list_pages(
self,
user: str,
tag: Optional[str] = None,
limit: int = 50
) -> WikiPageList:
"""
List pages for a user, optionally filtered by tag.
Args:
user: User identifier
tag: Optional tag filter (dossier)
limit: Maximum pages to return
Returns:
WikiPageList with pages and metadata
"""
namespace = self._get_user_namespace(user)
# Get pages with filtering
pages = await self.wiki_client.list_pages(
path_prefix=namespace,
tags=[tag] if tag else None,
limit=limit
)
# Convert to summary format
summaries = [
WikiPageSummary(
id=p["id"],
path=p["path"],
title=p["title"],
description=p.get("description"),
tags=p.get("tags", []),
updated_at=p.get("updatedAt"),
is_published=p.get("isPublished", True)
)
for p in pages
]
return WikiPageList(
pages=summaries,
total=len(summaries),
filtered_by_tag=tag,
user=user
)
async def get_page(self, page_id: int, user: str) -> Optional[WikiPage]:
"""
Get a single page by ID.
Args:
page_id: Page ID
user: User identifier (for validation)
Returns:
WikiPage or None if not found or access denied
Note: Validates that page belongs to user's namespace
"""
page = await self.wiki_client.get_page(page_id)
if not page:
return None
# Validate page is in user's namespace
namespace = self._get_user_namespace(user)
page_path = "/" + page["path"].lstrip("/") # Normalize path with leading slash
if not page_path.startswith(namespace):
logger.warning(f"User {user} attempted to access page outside namespace: {page['path']}")
return None
return WikiPage(
id=page["id"],
path=page["path"],
title=page["title"],
description=page.get("description"),
content=page.get("content"),
tags=page.get("tags", []),
created_at=page.get("createdAt"),
updated_at=page.get("updatedAt"),
is_published=page.get("isPublished", True),
editor=page.get("editor")
)
async def create_page(self, page_data: WikiPageCreate) -> WikiPage:
"""
Create a new wiki page.
Args:
page_data: Page creation data
Returns:
Created WikiPage
Raises:
ValueError: If creation fails
"""
user = page_data.user or DEFAULT_USER
# Ensure path is in user's namespace
full_path = self._ensure_user_path(page_data.path, user)
try:
created = await self.wiki_client.create_page(
path=full_path,
title=page_data.title,
content=page_data.content,
description=page_data.description or "",
tags=page_data.tags,
is_published=page_data.is_published,
editor=page_data.editor
)
# Fetch full page details
page = await self.wiki_client.get_page(created["id"])
if not page:
raise ValueError("Page created but could not be retrieved")
return WikiPage(
id=page["id"],
path=page["path"],
title=page["title"],
description=page.get("description"),
content=page.get("content"),
tags=page.get("tags", []),
created_at=page.get("createdAt"),
updated_at=page.get("updatedAt"),
is_published=page.get("isPublished", True),
editor=page.get("editor")
)
except Exception as e:
logger.error(f"Failed to create page: {e}", exc_info=True)
raise ValueError(f"Failed to create page: {str(e)}")
async def update_page(
self,
page_id: int,
page_data: WikiPageUpdate,
user: str
) -> WikiPage:
"""
Update an existing page.
Args:
page_id: Page ID to update
page_data: Update data
user: User identifier (for validation)
Returns:
Updated WikiPage
Raises:
ValueError: If page not found or update fails
"""
# Verify page exists and belongs to user
existing = await self.get_page(page_id, user)
if not existing:
raise ValueError(f"Page {page_id} not found or access denied")
try:
await self.wiki_client.update_page(
page_id=page_id,
content=page_data.content,
title=page_data.title,
description=page_data.description,
tags=page_data.tags,
is_published=True # Always keep pages published for internal wiki
)
# Fetch updated page
updated = await self.get_page(page_id, user)
if not updated:
raise ValueError("Page updated but could not be retrieved")
return updated
except Exception as e:
logger.error(f"Failed to update page {page_id}: {e}", exc_info=True)
raise ValueError(f"Failed to update page: {str(e)}")
async def delete_page(self, page_id: int, user: str) -> bool:
"""
Delete a page.
Args:
page_id: Page ID to delete
user: User identifier (for validation)
Returns:
True if deleted successfully
Raises:
ValueError: If page not found or deletion fails
"""
# Verify page exists and belongs to user
existing = await self.get_page(page_id, user)
if not existing:
raise ValueError(f"Page {page_id} not found or access denied")
try:
await self.wiki_client.delete_page(page_id)
logger.info(f"Deleted page {page_id} for user {user}")
return True
except Exception as e:
logger.error(f"Failed to delete page {page_id}: {e}", exc_info=True)
raise ValueError(f"Failed to delete page: {str(e)}")
async def search_pages(
self,
query: str,
user: str,
limit: int = 20
) -> List[WikiPageSummary]:
"""
Search pages in user's namespace.
Args:
query: Search query
user: User identifier
limit: Maximum results
Returns:
List of matching pages
"""
namespace = self._get_user_namespace(user)
results = await self.wiki_client.search_pages(
query=query,
path_prefix=namespace
)
# Convert to summaries (limit results)
return [
WikiPageSummary(
id=r["id"],
path=r["path"],
title=r["title"],
description=r.get("description"),
tags=[], # Search results don't include tags
updated_at=None,
is_published=True
)
for r in results[:limit]
]
async def move_page(
self,
page_id: int,
new_path: str,
user: str
) -> bool:
"""
Move/rename a page.
Args:
page_id: Page ID to move
new_path: New path (within user namespace)
user: User identifier
Returns:
True if moved successfully
Raises:
ValueError: If operation fails
"""
# Verify page exists and belongs to user
existing = await self.get_page(page_id, user)
if not existing:
raise ValueError(f"Page {page_id} not found or access denied")
# Ensure new path is in user's namespace
full_new_path = self._ensure_user_path(new_path, user)
try:
success = await self.wiki_client.move_page(page_id, full_new_path)
if success:
logger.info(f"Moved page {page_id} to {full_new_path}")
return success
except Exception as e:
logger.error(f"Failed to move page {page_id}: {e}", exc_info=True)
raise ValueError(f"Failed to move page: {str(e)}")
# Dossier operations (tag-based)
async def list_dossiers(self, user: str) -> DossierList:
"""
List all dossiers (unique tags) for a user.
Args:
user: User identifier
Returns:
DossierList with all dossiers
"""
# Get all pages for user
pages = await self.list_pages(user, limit=1000)
# Collect unique tags
tag_counts: Dict[str, int] = {}
for page in pages.pages:
for tag in page.tags:
tag_counts[tag] = tag_counts.get(tag, 0) + 1
# Create dossier info for each tag
dossiers = [
DossierInfo(
name=tag,
title=tag.replace("-", " ").title(),
description=f"Dossier for {tag}",
page_count=count,
index_page_id=None,
index_page_path=None,
created_at=None
)
for tag, count in tag_counts.items()
]
return DossierList(
dossiers=sorted(dossiers, key=lambda d: d.page_count, reverse=True),
total=len(dossiers),
user=user
)
async def get_dossier_pages(
self,
dossier_name: str,
user: str,
limit: int = 100
) -> WikiPageList:
"""
Get all pages in a dossier (by tag).
Args:
dossier_name: Dossier name (tag)
user: User identifier
limit: Maximum pages
Returns:
WikiPageList filtered by dossier tag
"""
return await self.list_pages(user, tag=dossier_name, limit=limit)