deerflow-code/offline-backend-20260512/backend/app/gateway/routers/external_llmwiki.py
2026-09-07 18:24:55 +08:00

134 lines
6.5 KiB
Python

"""API-key protected public Wiki vector search backed only by local data."""
from __future__ import annotations
import time
from collections import defaultdict, deque
from typing import Any
from fastapi import APIRouter, HTTPException, Request
from pydantic import BaseModel, Field
from deerflow.integrations.weknora.local_index.search import WikiIndexError
router = APIRouter(prefix="/api/external/llmwiki/wiki", tags=["external-llmwiki"])
_requests: dict[str, deque[float]] = defaultdict(deque)
class ExternalWikiSearchRequest(BaseModel):
query: str = Field(..., min_length=1, max_length=10000)
knowledge_base_ids: list[str] = Field(default_factory=list, max_length=100)
limit: int | None = Field(default=None, ge=1, le=100)
include_content: bool = False
def _config(request: Request):
config = request.app.state.config.llmwiki.local_wiki_index
if not config.enabled or not config.external_api.enabled:
raise HTTPException(status_code=404, detail="Not found")
return config, config.external_api
def _rate_limit(request: Request, rpm: int) -> None:
ip = request.client.host if request.client else "unknown"
key = f"{getattr(request.state, 'external_api_key_fingerprint', 'unknown')}:{ip}"
now = time.monotonic()
queue = _requests[key]
while queue and queue[0] <= now - 60:
queue.popleft()
if len(queue) >= rpm:
raise HTTPException(status_code=429, detail="Rate limit exceeded")
queue.append(now)
def _is_conversation_deposit(mapping: dict[str, Any]) -> bool:
return str(mapping.get("owner_user_id") or "") == "system" and str(mapping.get("name") or "").strip() == "对话沉淀"
async def _external_mappings(request: Request, requested: list[str]) -> list[dict[str, Any]]:
mappings = await request.app.state.llmwiki_store.list_visible("system", scope="all", is_admin=True)
allowed = {row["id"]: row for row in mappings if row.get("publication_status") == "published" and row.get("external_search_enabled", False) and row.get("wiki_index_enabled", True) and not _is_conversation_deposit(row)}
if requested:
ids = list(dict.fromkeys(requested))
if any(mapping_id not in allowed for mapping_id in ids):
raise HTTPException(status_code=404, detail="Knowledge base not found")
return [allowed[mapping_id] for mapping_id in ids]
return list(allowed.values())
@router.post("/vector-search")
async def external_vector_search(request: Request, body: ExternalWikiSearchRequest) -> dict[str, Any]:
_, external = _config(request)
_rate_limit(request, external.requests_per_minute)
if len(body.query) > external.max_query_chars:
raise HTTPException(status_code=413, detail="Query is too large")
if len(body.knowledge_base_ids) > external.max_knowledge_bases:
raise HTTPException(status_code=413, detail="Too many knowledge bases")
limit = body.limit or external.default_results
if limit > external.max_results:
raise HTTPException(status_code=400, detail="Result limit is too large")
mappings = await _external_mappings(request, body.knowledge_base_ids)
search = getattr(request.app.state, "llmwiki_vector_search", None)
if search is None:
raise HTTPException(status_code=503, detail={"code": "WIKI_INDEX_UNAVAILABLE"})
try:
result = await search.search(body.query, mappings, top_k_pages=limit, external=True)
except WikiIndexError as exc:
raise HTTPException(status_code=exc.status_code, detail={"code": exc.code, "message": str(exc)}) from None
# Re-read publication/external flags after vector computation so a switch
# turned off during an in-flight request cannot leak a cached result.
still_allowed = {row["id"] for row in await _external_mappings(request, [])}
result["results"] = [item for item in result["results"] if item.get("knowledge_base_id") in still_allowed]
remaining = external.include_content_max_chars if body.include_content else 0
items = []
for result_item in result["results"]:
section = (result_item.get("matched_sections") or [{}])[0]
snippet = str(section.get("content") or "")[:1000]
content = ""
if remaining > 0:
content = "\n\n".join(str(item.get("content") or "") for item in result_item.get("matched_sections") or [])[:remaining]
remaining -= len(content)
item = {
"knowledge_base_id": result_item["knowledge_base_id"],
"knowledge_base_name": result_item.get("knowledge_base_name") or "",
"wiki_slug": result_item["wiki_slug"],
"title": result_item.get("title") or "",
"summary": result_item.get("summary") or "",
"matched_heading": section.get("heading"),
"snippet": snippet,
"page_type": result_item.get("page_type") or "article",
"page_version": result_item.get("page_version") or 0,
"updated_at": result_item.get("updated_at"),
"score": result_item.get("score") or 0,
}
if body.include_content:
item["content"] = content
items.append(item)
return {"query": body.query, "mode": "vector", "items": items, "count": len(items), "partial": result.get("partial", False)}
@router.get("/pages/{mapping_id}/{slug:path}")
async def external_wiki_page(request: Request, mapping_id: str, slug: str) -> dict[str, Any]:
_, external = _config(request)
_rate_limit(request, external.requests_per_minute)
mappings = await _external_mappings(request, [mapping_id])
if not mappings:
raise HTTPException(status_code=404, detail="Wiki page not found")
page = await request.app.state.llmwiki_index_store.get_page(mapping_id, slug.strip("/"))
if page is None or page.get("status") != "published" or page.get("index_status") != "ready" or page.get("is_remote_deleted"):
raise HTTPException(status_code=404, detail="Wiki page not found")
# Recheck after reading the page so publication/external revocation takes
# effect even when it races this detail request.
await _external_mappings(request, [mapping_id])
return {
"knowledge_base_id": mapping_id,
"knowledge_base_name": mappings[0].get("name") or "",
"wiki_slug": page["slug"],
"title": page.get("title") or "",
"summary": page.get("summary") or "",
"content": str(page.get("content_md") or "")[: external.include_content_max_chars],
"page_type": page.get("page_type") or "article",
"page_version": page.get("page_version") or 0,
"updated_at": page.get("remote_updated_at") or page.get("updated_at"),
}