148 lines
5.2 KiB
Python
148 lines
5.2 KiB
Python
"""Authenticated proxy for WeKnora source-context previews.
|
|
|
|
The browser may ask for the source behind a DeerFlow skill citation, but it
|
|
must never receive the narrow WeKnora publication token. This router reads
|
|
the same server-side MCP configuration as the retrieval skill and forwards
|
|
only an authenticated, read-only context request.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from dataclasses import dataclass
|
|
from pathlib import Path
|
|
from urllib.parse import quote
|
|
|
|
import httpx
|
|
from fastapi import APIRouter, HTTPException, Query
|
|
|
|
from deerflow.config.extensions_config import get_extensions_config
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
router = APIRouter(prefix="/api/weknora", tags=["weknora"])
|
|
|
|
_MCP_SERVER_NAME = "weknora-retrieval"
|
|
_DEFAULT_TIMEOUT_SECONDS = 15.0
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class WeKnoraConnection:
|
|
base_url: str
|
|
share_token: str
|
|
timeout_seconds: float
|
|
|
|
|
|
def _read_secret_file(path_value: str) -> str:
|
|
try:
|
|
value = Path(path_value).read_text(encoding="utf-8").strip()
|
|
except OSError as exc:
|
|
raise RuntimeError("cannot read the WeKnora share-token file") from exc
|
|
if not value:
|
|
raise RuntimeError("the WeKnora share-token file is empty")
|
|
return value
|
|
|
|
|
|
def _connection_from_env(env: dict[str, str]) -> WeKnoraConnection:
|
|
base_url = env.get("WEKNORA_BASE_URL", "").strip().rstrip("/")
|
|
token_file = env.get("WEKNORA_DEERFLOW_SHARE_TOKEN_FILE", "").strip()
|
|
if not base_url:
|
|
raise RuntimeError("WEKNORA_BASE_URL is not configured")
|
|
if not token_file:
|
|
raise RuntimeError("WEKNORA_DEERFLOW_SHARE_TOKEN_FILE is not configured")
|
|
try:
|
|
timeout_seconds = float(
|
|
env.get("WEKNORA_REQUEST_TIMEOUT_SECONDS", _DEFAULT_TIMEOUT_SECONDS),
|
|
)
|
|
except (TypeError, ValueError) as exc:
|
|
raise RuntimeError("invalid WeKnora request timeout") from exc
|
|
if not 0 < timeout_seconds <= 60:
|
|
raise RuntimeError("invalid WeKnora request timeout")
|
|
return WeKnoraConnection(
|
|
base_url=base_url,
|
|
share_token=_read_secret_file(token_file),
|
|
timeout_seconds=timeout_seconds,
|
|
)
|
|
|
|
|
|
def load_weknora_connection() -> WeKnoraConnection:
|
|
config = get_extensions_config()
|
|
server = config.mcp_servers.get(_MCP_SERVER_NAME)
|
|
if server is None or not server.enabled:
|
|
raise RuntimeError("the WeKnora retrieval integration is not enabled")
|
|
return _connection_from_env(server.env)
|
|
|
|
|
|
async def fetch_source_context(
|
|
client: httpx.AsyncClient,
|
|
connection: WeKnoraConnection,
|
|
knowledge_base_id: str,
|
|
chunk_id: str,
|
|
) -> dict:
|
|
"""Fetch one context window and return only the WeKnora response data."""
|
|
path = (
|
|
"/api/v1/deerflow/knowledge-bases/"
|
|
f"{quote(knowledge_base_id, safe='')}/chunks/"
|
|
f"{quote(chunk_id, safe='')}/context"
|
|
)
|
|
response = await client.get(
|
|
f"{connection.base_url}{path}",
|
|
headers={"X-DeerFlow-Share-Token": connection.share_token},
|
|
)
|
|
if response.status_code == 404:
|
|
raise HTTPException(status_code=404, detail="The original source is no longer available")
|
|
if response.status_code >= 400:
|
|
raise RuntimeError(f"WeKnora source context returned HTTP {response.status_code}")
|
|
try:
|
|
payload = response.json()
|
|
except ValueError as exc:
|
|
raise RuntimeError("WeKnora source context returned invalid JSON") from exc
|
|
if not isinstance(payload, dict) or payload.get("success") is not True:
|
|
raise RuntimeError("WeKnora source context was rejected")
|
|
data = payload.get("data")
|
|
if not isinstance(data, dict) or not isinstance(data.get("chunks"), list):
|
|
raise RuntimeError("WeKnora source context returned an invalid payload")
|
|
return data
|
|
|
|
|
|
@router.get("/source-context", summary="Preview the original context of a WeKnora retrieval hit")
|
|
async def get_source_context(
|
|
knowledge_base_id: str = Query(..., min_length=1, max_length=128),
|
|
chunk_id: str = Query(..., min_length=1, max_length=128),
|
|
) -> dict:
|
|
"""Return nearby original chunks for a result cited by the WeKnora skill."""
|
|
try:
|
|
connection = load_weknora_connection()
|
|
except RuntimeError:
|
|
logger.exception("WeKnora source-context proxy is not configured")
|
|
raise HTTPException(
|
|
status_code=503,
|
|
detail="Original source preview is temporarily unavailable",
|
|
) from None
|
|
|
|
try:
|
|
timeout = httpx.Timeout(connection.timeout_seconds)
|
|
async with httpx.AsyncClient(timeout=timeout) as client:
|
|
data = await fetch_source_context(
|
|
client,
|
|
connection,
|
|
knowledge_base_id,
|
|
chunk_id,
|
|
)
|
|
except HTTPException:
|
|
raise
|
|
except httpx.HTTPError:
|
|
logger.warning("WeKnora source-context request failed", exc_info=True)
|
|
raise HTTPException(
|
|
status_code=502,
|
|
detail="Original source preview is temporarily unavailable",
|
|
) from None
|
|
except RuntimeError:
|
|
logger.warning("WeKnora source-context response was invalid", exc_info=True)
|
|
raise HTTPException(
|
|
status_code=502,
|
|
detail="Original source preview is temporarily unavailable",
|
|
) from None
|
|
|
|
return {"success": True, "data": data}
|