feat(phase3): add The Librarian agent with library-desk integration
Library-Desk API Client: - Async HTTP client with httpx for library-desk API - HybridRAG search (vector + graph + web) - Wiki operations (search, get, list, create, update) - Smart page creation with HybridRAG research - Semantic vector search and knowledge graph queries - Dossier browsing and health checks Librarian Tools (11 total): - Research: hybrid_search, search_wiki, get_wiki_page, semantic_search - Browse: list_dossiers, get_dossier_pages, explore_knowledge_graph - Graph: find_related_entities - Write: create_wiki_page, update_wiki_page, smart_create_wiki_page Agent: - PydanticAI agent with research assistant personality - System prompt with research and writing workflows - Streaming support via run_librarian_stream() Capability: - LIBRARIAN_CAPABILITY definition for Household Registry - Automatic registration on startup 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,685 @@
|
||||
"""
|
||||
HTTP client for the Library-Desk API.
|
||||
|
||||
Provides async methods for all relevant library-desk endpoints:
|
||||
- HybridRAG queries
|
||||
- Wiki operations
|
||||
- Vector search
|
||||
- Knowledge graph queries
|
||||
"""
|
||||
from typing import Any, Optional
|
||||
|
||||
import httpx
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
from src.core.config import config
|
||||
from src.core.logging_config import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
# ============================================================================
|
||||
# Response Models
|
||||
# ============================================================================
|
||||
|
||||
class WikiPage(BaseModel):
|
||||
"""Wiki page from library-desk."""
|
||||
id: int
|
||||
path: str
|
||||
title: str
|
||||
description: Optional[str] = None
|
||||
content: Optional[str] = None
|
||||
tags: list[str] = Field(default_factory=list)
|
||||
created_at: Optional[str] = None
|
||||
updated_at: Optional[str] = None
|
||||
|
||||
|
||||
class WikiSearchResult(BaseModel):
|
||||
"""Search result from wiki search."""
|
||||
id: int
|
||||
path: str
|
||||
title: str
|
||||
description: Optional[str] = None
|
||||
locale: Optional[str] = None
|
||||
|
||||
|
||||
class VectorSearchResult(BaseModel):
|
||||
"""Result from semantic vector search."""
|
||||
page_id: int
|
||||
page_path: str
|
||||
page_title: str
|
||||
chunk_text: str
|
||||
score: float
|
||||
chunk_index: int
|
||||
|
||||
|
||||
class HybridSearchResult(BaseModel):
|
||||
"""Result from HybridRAG search."""
|
||||
source: str # "vector", "graph", "web"
|
||||
title: str
|
||||
content: str
|
||||
url: Optional[str] = None
|
||||
score: float
|
||||
page_id: Optional[int] = None
|
||||
metadata: dict[str, Any] = Field(default_factory=dict)
|
||||
|
||||
|
||||
class HybridRAGResponse(BaseModel):
|
||||
"""Full response from HybridRAG query."""
|
||||
results: list[HybridSearchResult] = Field(default_factory=list)
|
||||
keywords: list[str] = Field(default_factory=list)
|
||||
synonyms: list[str] = Field(default_factory=list)
|
||||
related_dossiers: list[str] = Field(default_factory=list)
|
||||
formatted_context: str = ""
|
||||
search_id: Optional[str] = None
|
||||
timing: dict[str, float] = Field(default_factory=dict)
|
||||
|
||||
|
||||
class GraphNode(BaseModel):
|
||||
"""Node from knowledge graph."""
|
||||
id: str
|
||||
labels: list[str] = Field(default_factory=list)
|
||||
properties: dict[str, Any] = Field(default_factory=dict)
|
||||
|
||||
|
||||
class Dossier(BaseModel):
|
||||
"""A dossier (tag-based collection)."""
|
||||
name: str
|
||||
page_count: int
|
||||
|
||||
|
||||
class ResearchSummary(BaseModel):
|
||||
"""Summary of research performed during smart-create."""
|
||||
wiki_results: int = 0
|
||||
web_results: int = 0
|
||||
graph_entities: int = 0
|
||||
keywords_extracted: int = 0
|
||||
timing_ms: int = 0
|
||||
|
||||
|
||||
class EntityLinking(BaseModel):
|
||||
"""Entity linking results from smart-create."""
|
||||
forward_links: int = 0
|
||||
backward_links: int = 0
|
||||
pages_updated: int = 0
|
||||
|
||||
|
||||
class SmartCreateResponse(BaseModel):
|
||||
"""Response from smart-create wiki page endpoint."""
|
||||
page: WikiPage
|
||||
research_summary: ResearchSummary = Field(default_factory=ResearchSummary)
|
||||
sources_used: int = 0
|
||||
search_id: Optional[str] = None
|
||||
entity_linking: EntityLinking = Field(default_factory=EntityLinking)
|
||||
|
||||
|
||||
# ============================================================================
|
||||
# Client
|
||||
# ============================================================================
|
||||
|
||||
class LibraryDeskClient:
|
||||
"""
|
||||
Async HTTP client for Library-Desk API.
|
||||
|
||||
Usage:
|
||||
async with LibraryDeskClient() as client:
|
||||
results = await client.hybrid_search("docker kubernetes")
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
base_url: Optional[str] = None,
|
||||
api_key: Optional[str] = None,
|
||||
timeout: int = 60,
|
||||
):
|
||||
"""
|
||||
Initialize the client.
|
||||
|
||||
Args:
|
||||
base_url: Library-desk API URL (defaults to config)
|
||||
api_key: API key for authentication (defaults to config)
|
||||
timeout: Request timeout in seconds
|
||||
"""
|
||||
self.base_url = base_url or str(config.LIBRARY_DESK_HOST)
|
||||
self.api_key = api_key or config.LIBRARY_DESK_API_KEY
|
||||
self.timeout = timeout
|
||||
self._client: Optional[httpx.AsyncClient] = None
|
||||
|
||||
async def __aenter__(self) -> "LibraryDeskClient":
|
||||
"""Create HTTP client on context entry."""
|
||||
headers = {}
|
||||
if self.api_key:
|
||||
headers["Authorization"] = f"Bearer {self.api_key}"
|
||||
|
||||
self._client = httpx.AsyncClient(
|
||||
base_url=self.base_url,
|
||||
headers=headers,
|
||||
timeout=self.timeout,
|
||||
)
|
||||
return self
|
||||
|
||||
async def __aexit__(self, exc_type: Any, exc_val: Any, exc_tb: Any) -> None:
|
||||
"""Close HTTP client on context exit."""
|
||||
if self._client:
|
||||
await self._client.aclose()
|
||||
self._client = None
|
||||
|
||||
def _ensure_client(self) -> httpx.AsyncClient:
|
||||
"""Ensure client is initialized."""
|
||||
if self._client is None:
|
||||
raise RuntimeError(
|
||||
"Client not initialized. Use 'async with LibraryDeskClient() as client:'"
|
||||
)
|
||||
return self._client
|
||||
|
||||
# ========================================================================
|
||||
# HybridRAG
|
||||
# ========================================================================
|
||||
|
||||
async def hybrid_search(
|
||||
self,
|
||||
query: str,
|
||||
user: str = "jpmschweitzer",
|
||||
vector_limit: int = 10,
|
||||
graph_limit: int = 10,
|
||||
web_limit: int = 5,
|
||||
enable_reranking: bool = True,
|
||||
final_result_count: int = 10,
|
||||
) -> HybridRAGResponse:
|
||||
"""
|
||||
Execute HybridRAG search combining vector, graph, and web results.
|
||||
|
||||
Args:
|
||||
query: Search query
|
||||
user: User identifier for multi-tenancy
|
||||
vector_limit: Max results from vector search
|
||||
graph_limit: Max results from graph search
|
||||
web_limit: Max results from web search
|
||||
enable_reranking: Whether to rerank with LLM
|
||||
final_result_count: Number of final results after fusion
|
||||
|
||||
Returns:
|
||||
HybridRAGResponse with ranked results and context
|
||||
"""
|
||||
client = self._ensure_client()
|
||||
|
||||
payload = {
|
||||
"query": query,
|
||||
"config": {
|
||||
"vector_limit": vector_limit,
|
||||
"graph_limit": graph_limit,
|
||||
"web_limit": web_limit,
|
||||
"enable_reranking": enable_reranking,
|
||||
"final_result_count": final_result_count,
|
||||
},
|
||||
}
|
||||
|
||||
logger.info("library_desk_hybrid_search", query=query, user=user)
|
||||
|
||||
response = await client.post(
|
||||
"/query/hybrid",
|
||||
json=payload,
|
||||
params={"user": user},
|
||||
)
|
||||
response.raise_for_status()
|
||||
|
||||
data = response.json()
|
||||
|
||||
# Parse results
|
||||
results = []
|
||||
for r in data.get("results", []):
|
||||
results.append(HybridSearchResult(
|
||||
source=r.get("source", "unknown"),
|
||||
title=r.get("title", ""),
|
||||
content=r.get("content", ""),
|
||||
url=r.get("url"),
|
||||
score=r.get("score", 0.0),
|
||||
page_id=r.get("page_id"),
|
||||
metadata=r.get("metadata", {}),
|
||||
))
|
||||
|
||||
return HybridRAGResponse(
|
||||
results=results,
|
||||
keywords=data.get("keywords", []),
|
||||
synonyms=data.get("synonyms", []),
|
||||
related_dossiers=data.get("related_dossiers", []),
|
||||
formatted_context=data.get("formatted_context", ""),
|
||||
search_id=data.get("search_id"),
|
||||
timing=data.get("timing", {}),
|
||||
)
|
||||
|
||||
# ========================================================================
|
||||
# Wiki Operations
|
||||
# ========================================================================
|
||||
|
||||
async def search_wiki(
|
||||
self,
|
||||
query: str,
|
||||
user: str = "jpmschweitzer",
|
||||
limit: int = 20,
|
||||
) -> list[WikiSearchResult]:
|
||||
"""
|
||||
Search wiki pages by text.
|
||||
|
||||
Args:
|
||||
query: Search query
|
||||
user: User identifier
|
||||
limit: Maximum results
|
||||
|
||||
Returns:
|
||||
List of matching wiki pages
|
||||
"""
|
||||
client = self._ensure_client()
|
||||
|
||||
logger.debug("library_desk_wiki_search", query=query, user=user)
|
||||
|
||||
response = await client.get(
|
||||
"/wiki/search",
|
||||
params={"q": query, "user": user, "limit": limit},
|
||||
)
|
||||
response.raise_for_status()
|
||||
|
||||
data = response.json()
|
||||
return [WikiSearchResult(**r) for r in data.get("results", [])]
|
||||
|
||||
async def get_wiki_page(
|
||||
self,
|
||||
page_id: int,
|
||||
user: str = "jpmschweitzer",
|
||||
) -> WikiPage:
|
||||
"""
|
||||
Get a wiki page by ID.
|
||||
|
||||
Args:
|
||||
page_id: Page ID
|
||||
user: User identifier
|
||||
|
||||
Returns:
|
||||
WikiPage with full content
|
||||
"""
|
||||
client = self._ensure_client()
|
||||
|
||||
response = await client.get(
|
||||
f"/wiki/pages/{page_id}",
|
||||
params={"user": user},
|
||||
)
|
||||
response.raise_for_status()
|
||||
|
||||
return WikiPage(**response.json())
|
||||
|
||||
async def list_wiki_pages(
|
||||
self,
|
||||
user: str = "jpmschweitzer",
|
||||
tag: Optional[str] = None,
|
||||
limit: int = 50,
|
||||
) -> list[WikiPage]:
|
||||
"""
|
||||
List wiki pages, optionally filtered by tag.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
tag: Optional tag (dossier) to filter by
|
||||
limit: Maximum pages to return
|
||||
|
||||
Returns:
|
||||
List of wiki pages
|
||||
"""
|
||||
client = self._ensure_client()
|
||||
|
||||
params: dict[str, Any] = {"user": user, "limit": limit}
|
||||
if tag:
|
||||
params["tag"] = tag
|
||||
|
||||
response = await client.get("/wiki/pages", params=params)
|
||||
response.raise_for_status()
|
||||
|
||||
data = response.json()
|
||||
return [WikiPage(**p) for p in data.get("pages", [])]
|
||||
|
||||
async def create_wiki_page(
|
||||
self,
|
||||
title: str,
|
||||
path: str,
|
||||
content: str,
|
||||
user: str = "jpmschweitzer",
|
||||
description: str = "",
|
||||
tags: Optional[list[str]] = None,
|
||||
) -> WikiPage:
|
||||
"""
|
||||
Create a new wiki page.
|
||||
|
||||
Args:
|
||||
title: Page title
|
||||
path: Page path (e.g., "/projects/my-project")
|
||||
content: Markdown content
|
||||
user: User identifier
|
||||
description: Short description
|
||||
tags: List of tags (dossiers)
|
||||
|
||||
Returns:
|
||||
Created WikiPage
|
||||
"""
|
||||
client = self._ensure_client()
|
||||
|
||||
payload = {
|
||||
"title": title,
|
||||
"path": path,
|
||||
"content": content,
|
||||
"user": user,
|
||||
"description": description,
|
||||
"tags": tags or [],
|
||||
}
|
||||
|
||||
logger.info("library_desk_create_page", title=title, path=path)
|
||||
|
||||
response = await client.post("/wiki/pages", json=payload)
|
||||
response.raise_for_status()
|
||||
|
||||
return WikiPage(**response.json())
|
||||
|
||||
async def update_wiki_page(
|
||||
self,
|
||||
page_id: int,
|
||||
user: str = "jpmschweitzer",
|
||||
content: Optional[str] = None,
|
||||
title: Optional[str] = None,
|
||||
tags: Optional[list[str]] = None,
|
||||
description: Optional[str] = None,
|
||||
) -> WikiPage:
|
||||
"""
|
||||
Update an existing wiki page.
|
||||
|
||||
Supports partial updates - only provided fields are updated.
|
||||
Automatically triggers vector re-indexing and graph extraction.
|
||||
|
||||
Args:
|
||||
page_id: ID of the page to update
|
||||
user: User identifier
|
||||
content: New content (optional)
|
||||
title: New title (optional)
|
||||
tags: New tags list (optional)
|
||||
description: New description (optional)
|
||||
|
||||
Returns:
|
||||
Updated WikiPage
|
||||
"""
|
||||
client = self._ensure_client()
|
||||
|
||||
# Build update payload with only provided fields
|
||||
update_data: dict[str, Any] = {}
|
||||
if content is not None:
|
||||
update_data["content"] = content
|
||||
if title is not None:
|
||||
update_data["title"] = title
|
||||
if tags is not None:
|
||||
update_data["tags"] = tags
|
||||
if description is not None:
|
||||
update_data["description"] = description
|
||||
|
||||
logger.info(
|
||||
"library_desk_update_page",
|
||||
page_id=page_id,
|
||||
fields=list(update_data.keys()),
|
||||
)
|
||||
|
||||
response = await client.put(
|
||||
f"/wiki/pages/{page_id}",
|
||||
params={"user": user},
|
||||
json=update_data,
|
||||
)
|
||||
response.raise_for_status()
|
||||
|
||||
return WikiPage(**response.json())
|
||||
|
||||
async def smart_create_wiki_page(
|
||||
self,
|
||||
topic: str,
|
||||
tags: list[str],
|
||||
user: str = "jpmschweitzer",
|
||||
path: Optional[str] = None,
|
||||
include_web_research: bool = True,
|
||||
include_wiki_search: bool = True,
|
||||
) -> SmartCreateResponse:
|
||||
"""
|
||||
Create a wiki page with HybridRAG research.
|
||||
|
||||
This endpoint:
|
||||
1. Searches existing wiki, knowledge graph, and web for context
|
||||
2. Uses LLM to synthesize findings into structured content
|
||||
3. Creates the page with proper attribution
|
||||
4. Automatically links entities bidirectionally
|
||||
|
||||
Args:
|
||||
topic: The topic to research and create a page about
|
||||
tags: List of tags (dossiers) for the page
|
||||
user: User identifier
|
||||
path: Optional custom path (auto-generated from topic if not provided)
|
||||
include_web_research: Whether to include web search results
|
||||
include_wiki_search: Whether to include existing wiki content
|
||||
|
||||
Returns:
|
||||
SmartCreateResponse with page and research metadata
|
||||
"""
|
||||
client = self._ensure_client()
|
||||
|
||||
payload: dict[str, Any] = {
|
||||
"topic": topic,
|
||||
"tags": tags,
|
||||
"user": user,
|
||||
"include_web_research": include_web_research,
|
||||
"include_wiki_search": include_wiki_search,
|
||||
}
|
||||
if path is not None:
|
||||
payload["path"] = path
|
||||
|
||||
logger.info(
|
||||
"library_desk_smart_create",
|
||||
topic=topic,
|
||||
tags=tags,
|
||||
include_web=include_web_research,
|
||||
)
|
||||
|
||||
response = await client.post("/wiki/pages/smart-create", json=payload)
|
||||
response.raise_for_status()
|
||||
|
||||
data = response.json()
|
||||
|
||||
# Parse nested response
|
||||
page = WikiPage(**data.get("page", {}))
|
||||
research_summary = ResearchSummary(**data.get("research_summary", {}))
|
||||
entity_linking = EntityLinking(**data.get("entity_linking", {}))
|
||||
|
||||
return SmartCreateResponse(
|
||||
page=page,
|
||||
research_summary=research_summary,
|
||||
sources_used=data.get("sources_used", 0),
|
||||
search_id=data.get("search_id"),
|
||||
entity_linking=entity_linking,
|
||||
)
|
||||
|
||||
async def list_dossiers(
|
||||
self,
|
||||
user: str = "jpmschweitzer",
|
||||
) -> list[Dossier]:
|
||||
"""
|
||||
List all dossiers (tag collections) for a user.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
|
||||
Returns:
|
||||
List of dossiers with page counts
|
||||
"""
|
||||
client = self._ensure_client()
|
||||
|
||||
response = await client.get(
|
||||
"/wiki/dossiers",
|
||||
params={"user": user},
|
||||
)
|
||||
response.raise_for_status()
|
||||
|
||||
data = response.json()
|
||||
return [Dossier(**d) for d in data.get("dossiers", [])]
|
||||
|
||||
# ========================================================================
|
||||
# Vector Search
|
||||
# ========================================================================
|
||||
|
||||
async def semantic_search(
|
||||
self,
|
||||
query: str,
|
||||
user: str = "jpmschweitzer",
|
||||
limit: int = 10,
|
||||
score_threshold: float = 0.5,
|
||||
) -> list[VectorSearchResult]:
|
||||
"""
|
||||
Perform semantic (vector) search over documents.
|
||||
|
||||
Args:
|
||||
query: Natural language query
|
||||
user: User identifier
|
||||
limit: Maximum results
|
||||
score_threshold: Minimum similarity score
|
||||
|
||||
Returns:
|
||||
List of matching document chunks with scores
|
||||
"""
|
||||
client = self._ensure_client()
|
||||
|
||||
payload = {
|
||||
"query": query,
|
||||
"user": user,
|
||||
"limit": limit,
|
||||
"score_threshold": score_threshold,
|
||||
}
|
||||
|
||||
logger.debug("library_desk_semantic_search", query=query)
|
||||
|
||||
response = await client.post("/vector/search", json=payload)
|
||||
response.raise_for_status()
|
||||
|
||||
data = response.json()
|
||||
return [VectorSearchResult(**r) for r in data.get("results", [])]
|
||||
|
||||
# ========================================================================
|
||||
# Knowledge Graph
|
||||
# ========================================================================
|
||||
|
||||
async def query_graph(
|
||||
self,
|
||||
cypher_query: str,
|
||||
user: str = "jpmschweitzer",
|
||||
parameters: Optional[dict[str, Any]] = None,
|
||||
) -> list[dict[str, Any]]:
|
||||
"""
|
||||
Execute a Cypher query on the knowledge graph.
|
||||
|
||||
Note: Query is automatically scoped to user's data.
|
||||
|
||||
Args:
|
||||
cypher_query: Cypher query string
|
||||
user: User identifier
|
||||
parameters: Query parameters
|
||||
|
||||
Returns:
|
||||
List of result records
|
||||
"""
|
||||
client = self._ensure_client()
|
||||
|
||||
payload = {
|
||||
"query": cypher_query,
|
||||
"user": user,
|
||||
"parameters": parameters or {},
|
||||
}
|
||||
|
||||
logger.debug("library_desk_graph_query", query=cypher_query[:100])
|
||||
|
||||
response = await client.post("/graph/query", json=payload)
|
||||
response.raise_for_status()
|
||||
|
||||
return response.json().get("records", [])
|
||||
|
||||
async def list_graph_nodes(
|
||||
self,
|
||||
user: str = "jpmschweitzer",
|
||||
node_type: Optional[str] = None,
|
||||
limit: int = 100,
|
||||
) -> list[GraphNode]:
|
||||
"""
|
||||
List nodes in the knowledge graph.
|
||||
|
||||
Args:
|
||||
user: User identifier
|
||||
node_type: Optional filter by type (Document, Person, Concept, etc.)
|
||||
limit: Maximum nodes
|
||||
|
||||
Returns:
|
||||
List of graph nodes
|
||||
"""
|
||||
client = self._ensure_client()
|
||||
|
||||
params: dict[str, Any] = {"user": user, "limit": limit}
|
||||
if node_type:
|
||||
params["node_type"] = node_type
|
||||
|
||||
response = await client.get("/graph/nodes", params=params)
|
||||
response.raise_for_status()
|
||||
|
||||
data = response.json()
|
||||
return [GraphNode(**n) for n in data.get("nodes", [])]
|
||||
|
||||
async def get_graph_node(
|
||||
self,
|
||||
node_id: str,
|
||||
user: str = "jpmschweitzer",
|
||||
) -> dict[str, Any]:
|
||||
"""
|
||||
Get detailed information about a graph node.
|
||||
|
||||
Args:
|
||||
node_id: Node ID
|
||||
user: User identifier
|
||||
|
||||
Returns:
|
||||
Node with relationships and connected nodes
|
||||
"""
|
||||
client = self._ensure_client()
|
||||
|
||||
response = await client.get(
|
||||
f"/graph/nodes/{node_id}",
|
||||
params={"user": user},
|
||||
)
|
||||
response.raise_for_status()
|
||||
|
||||
return response.json()
|
||||
|
||||
# ========================================================================
|
||||
# Health Check
|
||||
# ========================================================================
|
||||
|
||||
async def health_check(self) -> bool:
|
||||
"""
|
||||
Check if library-desk is healthy.
|
||||
|
||||
Returns:
|
||||
True if healthy, False otherwise
|
||||
"""
|
||||
try:
|
||||
client = self._ensure_client()
|
||||
response = await client.get("/health")
|
||||
return response.status_code == 200
|
||||
except Exception as e:
|
||||
logger.warning("library_desk_health_check_failed", error=str(e))
|
||||
return False
|
||||
|
||||
|
||||
# Global client factory
|
||||
async def get_library_client() -> LibraryDeskClient:
|
||||
"""
|
||||
Get a library-desk client instance.
|
||||
|
||||
Usage:
|
||||
async with get_library_client() as client:
|
||||
results = await client.hybrid_search("query")
|
||||
"""
|
||||
return LibraryDeskClient()
|
||||
Reference in New Issue
Block a user