"""cmzs LLMWiki BFF backed by an independently deployed WeKnora.""" from __future__ import annotations import asyncio import hashlib import json import logging import mimetypes import re import time from dataclasses import replace from datetime import UTC, datetime from io import BytesIO from typing import Any, Literal from urllib.parse import quote, unquote, urlencode, urlsplit import httpx from fastapi import APIRouter, BackgroundTasks, File, HTTPException, Query, Request, UploadFile from fastapi.responses import JSONResponse, RedirectResponse, Response, StreamingResponse from pydantic import BaseModel, Field from app.gateway.csrf_middleware import is_secure_request from app.gateway.deps import get_agent_store, get_optional_user_from_request, get_thread_store from app.gateway.llmwiki_deposit import ( CONVERSATION_DEPOSIT_KB_DESCRIPTION, CONVERSATION_DEPOSIT_KB_NAME, CONVERSATION_DEPOSIT_OWNER_USER_ID, ConversationTurnDeposit, deposit_conversation_turn, ensure_conversation_deposit_mapping, ) from app.gateway.weknora_embed import ( WEKNORA_EMBED_API_PREFIX as _WEKNORA_EMBED_API_PREFIX, ) from app.gateway.weknora_embed import ( WEKNORA_EMBED_COOKIE as _WEKNORA_EMBED_COOKIE, ) from app.gateway.weknora_embed import ( WEKNORA_EMBED_TTL_SECONDS as _WEKNORA_EMBED_TTL_SECONDS, ) from app.gateway.weknora_embed import ( WeKnoraEmbedContext, ) from app.gateway.weknora_embed import ( sign_weknora_embed_payload as _sign_embed_payload, ) from app.gateway.weknora_embed import ( verify_weknora_embed_cookie as _verify_embed_cookie, ) from app.gateway.weknora_embed import ( verify_weknora_signed_payload as _verify_signed_payload, ) from deerflow.config.agents_config import load_agent_config from deerflow.config.llmwiki_config import LlmWikiRuntimeOverride, normalize_service_url from deerflow.config.system_settings import load_system_settings, save_system_settings from deerflow.integrations.weknora.client import WeKnoraClient, WeKnoraError from deerflow.integrations.weknora.local_index.retrieval import is_fully_vectorized from deerflow.integrations.weknora.local_index.search import WikiIndexError from deerflow.integrations.weknora.runtime import build_weknora_client, get_resolved_llmwiki_runtime from deerflow.persistence.llmwiki import LlmWikiStore logger = logging.getLogger(__name__) router = APIRouter(tags=["llmwiki"]) proxy_router = APIRouter(tags=["llmwiki-proxy"]) _MAX_UPLOAD_BYTES = 100 * 1024 * 1024 _WIKI_FILE_ACCESS_TTL_SECONDS = 24 * 60 * 60 _PROTECTED_WIKI_FILE_PATH = re.compile( r"^(?:storage://[0-9A-Za-z_-]+/)?(?:resource|local|minio|cos|tos|s3|oss|ks3|obs)://", re.IGNORECASE, ) _WEKNORA_ADMIN_SESSION_CACHE: dict[str, dict[str, Any]] = {} _WEKNORA_PUBLIC_PREFIX = "/deerflow" _WEKNORA_PUBLIC_ROOT_PATHS = ( "/platform", "/assets", "/css", "/fonts", "/img", "/images", "/js", "/locales", "/media", "/static", "/tdesign-icons", "/api/v1", "/favicon.ico", "/favicon.svg", "/config.js", "/logo.svg", "/manifest.webmanifest", ) _WEKNORA_PUBLIC_RELATIVE_ROOT_PATHS = ( "assets", "css", "fonts", "img", "images", "js", "locales", "media", "static", "tdesign-icons", ) _WEKNORA_PUBLIC_RELATIVE_FILE_PATHS = ( "config.js", "favicon.ico", "favicon.svg", "logo.svg", "manifest.webmanifest", ) _WEKNORA_STATIC_FILE_SUFFIXES = ( ".avif", ".css", ".gif", ".ico", ".jpeg", ".jpg", ".js", ".json", ".map", ".mjs", ".png", ".svg", ".ttf", ".txt", ".webmanifest", ".webp", ".woff", ".woff2", ) _WEKNORA_STATIC_ROOT_PREFIXES = ( "/assets/", "/css/", "/fonts/", "/img/", "/images/", "/js/", "/locales/", "/media/", "/static/", "/tdesign-icons/", ) _WEKNORA_ROOT_ASSET_FILES = ( "/config.js", "/favicon.ico", "/favicon.svg", "/logo.svg", "/manifest.webmanifest", ) _HOP_BY_HOP_RESPONSE_HEADERS = { "connection", "content-encoding", "content-length", "keep-alive", "proxy-authenticate", "proxy-authorization", "te", "trailer", "transfer-encoding", "upgrade", "x-frame-options", "content-security-policy", } class LlmWikiSettingsUpdate(BaseModel): override_enabled: bool = True api_base_url: str = "" web_base_url: str = "" class LlmWikiConnectionTest(BaseModel): api_base_url: str | None = None class KnowledgeBaseCreate(BaseModel): name: str = Field(..., min_length=1, max_length=255) description: str = Field(default="", max_length=2048) type: Literal["document", "faq"] = "document" wiki_enabled: bool = False class KnowledgeBaseUpdate(BaseModel): name: str | None = Field(default=None, min_length=1, max_length=255) description: str | None = Field(default=None, max_length=2048) class SearchRequest(BaseModel): query: str = Field(..., min_length=1, max_length=4000) scope: Literal["personal", "public", "all"] = "all" knowledge_base_ids: list[str] = Field(default_factory=list, max_length=100) top_k: int = Field(default=8, ge=1, le=50) include_drafts: bool = False class KnowledgeUrlImport(BaseModel): url: str = Field(..., min_length=1, max_length=4096) class ManualKnowledgeCreate(BaseModel): title: str = Field(..., min_length=1, max_length=255) content: str = Field(..., min_length=1, max_length=2_000_000) class WikiPageCreate(BaseModel): slug: str = Field(..., min_length=1, max_length=512) title: str = Field(..., min_length=1, max_length=255) content: str = Field(..., min_length=1, max_length=2_000_000) summary: str | None = Field(default=None, max_length=200_000) page_type: str | None = Field(default="article", min_length=1, max_length=64) status: str | None = Field(default="published", min_length=1, max_length=64) aliases: list[str] | None = Field(default=None, max_length=100) parent_slug: str | None = Field(default=None, max_length=512) category_path: list[str] | None = Field(default=None, max_length=50) folder_id: str | None = Field(default=None, max_length=128) wiki_path: str | None = Field(default=None, max_length=1024) class WikiPageUpdate(BaseModel): title: str | None = Field(default=None, min_length=1, max_length=255) content: str | None = Field(default=None, min_length=1, max_length=2_000_000) summary: str | None = Field(default=None, max_length=200_000) page_type: str | None = Field(default=None, min_length=1, max_length=64) status: str | None = Field(default=None, min_length=1, max_length=64) aliases: list[str] | None = Field(default=None, max_length=100) parent_slug: str | None = Field(default=None, max_length=512) category_path: list[str] | None = Field(default=None, max_length=50) folder_id: str | None = Field(default=None, max_length=128) wiki_path: str | None = Field(default=None, max_length=1024) class ConversationDepositCreate(BaseModel): thread_id: str = Field(..., min_length=1, max_length=128) question: str = Field(default="", max_length=200_000) answer: str = Field(default="", max_length=2_000_000) human_message_id: str | None = Field(default=None, max_length=128) assistant_message_id: str = Field(..., min_length=1, max_length=128) assistant_id: str | None = Field(default=None, max_length=128) knowledge_base_ids: list[str] = Field(default_factory=list, max_length=100) def _store(request: Request) -> LlmWikiStore: value = getattr(request.app.state, "llmwiki_store", None) if value is None: raise HTTPException(status_code=503, detail="LLMWiki metadata store is unavailable") return value def _schedule_local_wiki_sync(request: Request, mapping: dict[str, Any]) -> None: service = getattr(request.app.state, "llmwiki_sync_service", None) config = request.app.state.config.llmwiki.local_wiki_index if service is not None and config.enabled and config.auto_sync and mapping.get("wiki_index_enabled", True): from app.gateway.llmwiki_index_scheduler import spawn_local_wiki_index_job # Event-driven refresh is incremental. Content hash + embedding # fingerprint decide which pages need encoding; unchanged Wiki pages # must retain and reuse their existing vectors. spawn_local_wiki_index_job(request.app, service.sync_mapping(mapping, force=False)) async def _actor(request: Request) -> tuple[str, bool]: user = await get_optional_user_from_request(request) if user is None: return "default", True return str(user.id), getattr(user, "system_role", None) == "admin" async def _require_admin(request: Request) -> tuple[str, bool]: actor = await _actor(request) if not actor[1]: raise HTTPException(status_code=403, detail="Only administrators can configure WeKnora") return actor def _runtime_or_legacy() -> dict[str, Any]: runtime = get_resolved_llmwiki_runtime() return { "provider": runtime.provider, "api_base_url": runtime.api_base_url, "web_base_url": runtime.web_base_url, "override_enabled": runtime.override_enabled, } def _client_or_503() -> WeKnoraClient: runtime = get_resolved_llmwiki_runtime() if not runtime.weknora_enabled: raise HTTPException(status_code=409, detail="WeKnora mode is not enabled") try: return build_weknora_client(runtime) except RuntimeError as exc: raise HTTPException(status_code=503, detail=str(exc)) from None def _raise_upstream(exc: WeKnoraError) -> None: status = 422 if exc.status_code == 400 else exc.status_code if exc.status_code in {404, 502, 504} else 502 raise HTTPException(status_code=status, detail=str(exc)) from None def _json_response(payload: Any, *, status_code: int = 200) -> Response: return Response( content=json.dumps(payload, ensure_ascii=False, separators=(",", ":")).encode("utf-8"), status_code=status_code, media_type="application/json", ) def _runtime_for_iframe_proxy(request: Request | None = None): runtime = get_resolved_llmwiki_runtime(getattr(request.app.state, "config", None) if request is not None else None) if not runtime.weknora_enabled: raise HTTPException(status_code=409, detail="WeKnora mode is not enabled") if not runtime.web_base_url: raise HTTPException(status_code=422, detail="WeKnora web_base_url is required for iframe embedding") if not runtime.admin_email or not runtime.admin_password: raise HTTPException( status_code=422, detail="Configure llmwiki.weknora.admin_email and admin_password for WeKnora iframe SSO", ) return runtime def _admin_session_cache_key(runtime: Any) -> str: digest = hashlib.sha256(runtime.admin_password.encode("utf-8")).hexdigest()[:12] return f"{runtime.api_base_url}|{runtime.admin_email}|{digest}" def _extract_weknora_token(payload: dict[str, Any]) -> tuple[str, str, dict[str, Any] | None, dict[str, Any] | None]: source = payload.get("data") if isinstance(payload.get("data"), dict) else payload token = str(source.get("token") or source.get("access_token") or payload.get("token") or payload.get("access_token") or "") refresh_token = str(source.get("refresh_token") or source.get("refreshToken") or payload.get("refresh_token") or payload.get("refreshToken") or "") user = source.get("user") if isinstance(source.get("user"), dict) else payload.get("user") tenant = source.get("tenant") if isinstance(source.get("tenant"), dict) else source.get("active_tenant") if isinstance(source.get("active_tenant"), dict) else payload.get("tenant") return token, refresh_token, user if isinstance(user, dict) else None, tenant if isinstance(tenant, dict) else None async def _get_weknora_admin_session(runtime: Any, *, force_refresh: bool = False) -> dict[str, Any]: key = _admin_session_cache_key(runtime) cached = _WEKNORA_ADMIN_SESSION_CACHE.get(key) now = time.time() if cached and not force_refresh and float(cached.get("expires_at") or 0) > now + 60: return cached try: async with httpx.AsyncClient(timeout=httpx.Timeout(20.0), follow_redirects=False, trust_env=False) as client: response = await client.post( f"{runtime.api_base_url}/api/v1/auth/login", json={"email": runtime.admin_email, "password": runtime.admin_password}, headers={"Accept": "application/json"}, ) except httpx.TimeoutException as exc: logger.warning("WeKnora admin login timed out: base_url=%s", runtime.api_base_url) raise HTTPException(status_code=504, detail="WeKnora admin login timed out") from exc except httpx.HTTPError as exc: logger.warning("WeKnora admin login unavailable: base_url=%s error=%r", runtime.api_base_url, exc) raise HTTPException(status_code=502, detail="WeKnora admin login is unavailable") from exc if response.status_code >= 400: raise HTTPException(status_code=502, detail=f"WeKnora admin login failed ({response.status_code})") try: payload = response.json() except ValueError as exc: raise HTTPException(status_code=502, detail="WeKnora admin login returned invalid JSON") from exc if isinstance(payload, dict) and payload.get("success") is False: raise HTTPException(status_code=502, detail=str(payload.get("message") or "WeKnora admin login failed")) token, refresh_token, user, tenant = _extract_weknora_token(payload if isinstance(payload, dict) else {}) if not token: raise HTTPException(status_code=502, detail="WeKnora admin login did not return a token") session = { "token": token, "refresh_token": refresh_token, "user": user or {}, "tenant": tenant or {}, # WeKnora tokens do not consistently expose expiry in the login body. # Relogin periodically and retry once on upstream 401. "expires_at": now + 25 * 60, } _WEKNORA_ADMIN_SESSION_CACHE[key] = session return session def _tenant_id_from_session(session: dict[str, Any]) -> str: tenant = session.get("tenant") if isinstance(session.get("tenant"), dict) else {} user = session.get("user") if isinstance(session.get("user"), dict) else {} return str(tenant.get("id") or user.get("tenant_id") or "") async def _read_embed_context(request: Request) -> WeKnoraEmbedContext: ctx = _verify_embed_cookie(request.cookies.get(_WEKNORA_EMBED_COOKIE)) if not ctx.mapping_id or not ctx.weknora_id or not ctx.user_id: raise HTTPException(status_code=401, detail="WeKnora iframe session is invalid") store = _store(request) readable = await store.get_authorized(ctx.mapping_id, ctx.user_id, write=False, is_admin=ctx.is_admin) if readable is None or str(readable.get("weknora_id") or "") != ctx.weknora_id: raise HTTPException(status_code=403, detail="WeKnora iframe session is no longer authorized") writable = None if not _is_conversation_deposit_mapping(readable): writable = await store.get_authorized(ctx.mapping_id, ctx.user_id, write=True, is_admin=ctx.is_admin) return replace( ctx, can_write=writable is not None, is_conversation_deposit=_is_conversation_deposit_mapping(readable), allowed_weknora_ids=_embed_allowed_weknora_ids(ctx), ) def _embed_allowed_weknora_ids(ctx: WeKnoraEmbedContext) -> tuple[str, ...]: """Return the signed knowledge-base allowlist, always retaining the active base.""" return tuple(dict.fromkeys(value for value in (ctx.weknora_id, *ctx.allowed_weknora_ids) if value)) def _embed_context_payload(ctx: WeKnoraEmbedContext) -> dict[str, Any]: return { "mapping_id": ctx.mapping_id, "weknora_id": ctx.weknora_id, "user_id": ctx.user_id, "is_admin": ctx.is_admin, "can_write": ctx.can_write, "is_conversation_deposit": ctx.is_conversation_deposit, "allowed_weknora_ids": list(_embed_allowed_weknora_ids(ctx)), "exp": ctx.exp, } def _weknora_id_from_platform_path(path: str) -> str | None: prefix = "/platform/knowledge-bases/" if not path.startswith(prefix): return None raw_id = path.removeprefix(prefix) if not raw_id or "/" in raw_id: return None return unquote(raw_id) async def _switch_embed_context( request: Request, ctx: WeKnoraEmbedContext, weknora_id: str, ) -> WeKnoraEmbedContext: """Resolve a permitted native-UI switch to the exact cmzs mapping. The embedded WeKnora UI uses a service-admin token, so every target remains constrained to an ID signed into this short-lived session and is checked again against the cmzs visibility rules before its page can be proxied. """ if weknora_id not in _embed_allowed_weknora_ids(ctx): raise HTTPException(status_code=403, detail="Knowledge base is outside the embedded session scope") store = _store(request) mapped = await store.get_by_weknora_id(weknora_id) if mapped is None: raise HTTPException(status_code=403, detail="Knowledge base is not managed by cmzs") readable = await store.get_authorized( str(mapped["id"]), ctx.user_id, write=False, is_admin=ctx.is_admin, ) if readable is None or str(readable.get("weknora_id") or "") != weknora_id: raise HTTPException(status_code=403, detail="Knowledge base is no longer authorized") writable = None if not _is_conversation_deposit_mapping(readable): writable = await store.get_authorized( str(readable["id"]), ctx.user_id, write=True, is_admin=ctx.is_admin, ) return replace( ctx, mapping_id=str(readable["id"]), weknora_id=weknora_id, can_write=writable is not None, is_conversation_deposit=_is_conversation_deposit_mapping(readable), ) def _encode_upstream_path(path: str) -> str: if not path.startswith("/"): path = f"/{path}" return "/".join(quote(part, safe="") for part in path.split("/")) def _safe_response_headers(headers: httpx.Headers) -> dict[str, str]: result: dict[str, str] = {} for key, value in headers.items(): if key.lower() in _HOP_BY_HOP_RESPONSE_HEADERS: continue if key.lower() == "set-cookie": continue result[key] = value return result def _should_follow_weknora_binary_redirect(path: str) -> bool: """Keep protected preview/image bytes inside the DeerFlow proxy.""" segments = _path_segments(path) return (len(segments) == 5 and segments[:3] == ["api", "v1", "knowledge"] and segments[4] == "preview") or (len(segments) == 5 and segments[:3] == ["api", "v1", "knowledge-bases"] and segments[4] == "files") def _forward_request_headers(request: Request) -> dict[str, str]: blocked = { "host", "connection", "content-length", "cookie", "authorization", "x-forwarded-host", "x-forwarded-for", "x-forwarded-proto", } headers = {key: value for key, value in request.headers.items() if key.lower() not in blocked} headers.setdefault("Accept", "*/*") return headers def _path_segments(path: str) -> list[str]: return [segment for segment in path.strip("/").split("/") if segment] def _json_from_body(body: bytes, content_type: str | None) -> Any: if not body or "application/json" not in (content_type or "").lower(): return None try: return json.loads(body) except ValueError: return None def _body_kb_id(body: Any) -> str: return str(body.get("kb_id") or body.get("knowledge_base_id") or body.get("source_kb_id") or "") if isinstance(body, dict) else "" def _request_kb_id(request: Request, body: Any) -> str: return _body_kb_id(body) or str(request.query_params.get("kb_id") or request.query_params.get("knowledge_base_id") or request.query_params.get("source_kb_id") or "") async def _admin_json(runtime: Any, path: str, *, token: str, params: dict[str, Any] | None = None) -> dict[str, Any]: try: async with httpx.AsyncClient(timeout=httpx.Timeout(20.0), follow_redirects=False, trust_env=False) as client: response = await client.get( f"{runtime.api_base_url}{_encode_upstream_path(path)}", params=params, headers={"Accept": "application/json", "Authorization": f"Bearer {token}"}, ) except httpx.HTTPError as exc: logger.warning("WeKnora verification request failed: base_url=%s path=%s error=%r", runtime.api_base_url, path, exc) raise HTTPException(status_code=502, detail="WeKnora verification request failed") from exc if response.status_code == 404: raise HTTPException(status_code=404, detail="WeKnora resource was not found") if response.status_code >= 400: raise HTTPException(status_code=502, detail=f"WeKnora verification failed ({response.status_code})") try: payload = response.json() except ValueError as exc: raise HTTPException(status_code=502, detail="WeKnora verification returned invalid JSON") from exc return payload if isinstance(payload, dict) else {} def _payload_data(payload: dict[str, Any]) -> Any: return payload.get("data", payload) async def _verify_proxy_document_belongs_to_kb(runtime: Any, ctx: WeKnoraEmbedContext, document_id: str, token: str) -> None: payload = await _admin_json(runtime, f"/api/v1/knowledge/{document_id}", token=token) data = _payload_data(payload) if not isinstance(data, dict) or str(data.get("knowledge_base_id") or "") != ctx.weknora_id: raise HTTPException(status_code=403, detail="Document is outside the authorized WeKnora knowledge base") async def _verify_proxy_chunk_belongs_to_kb(runtime: Any, ctx: WeKnoraEmbedContext, chunk_id: str, token: str) -> None: payload = await _admin_json(runtime, f"/api/v1/chunks/by-id/{chunk_id}", token=token) data = _payload_data(payload) if not isinstance(data, dict) or str(data.get("knowledge_base_id") or "") != ctx.weknora_id: raise HTTPException(status_code=403, detail="Chunk is outside the authorized WeKnora knowledge base") async def _authorize_weknora_api_proxy( request: Request, ctx: WeKnoraEmbedContext, runtime: Any, path: str, body: bytes, token: str, ) -> None: method = request.method.upper() if method == "OPTIONS": return read_method = method in {"GET", "HEAD"} write_method = method in {"POST", "PUT", "PATCH", "DELETE"} if write_method and not ctx.can_write: raise HTTPException(status_code=403, detail="Current cmzs user cannot write this knowledge base") segments = _path_segments(path) if len(segments) < 2 or segments[0] != "api" or segments[1] != "v1": raise HTTPException(status_code=403, detail="Unsupported WeKnora API path") tail = segments[2:] json_body = _json_from_body(body, request.headers.get("content-type")) if tail[:1] == ["auth"]: if len(tail) >= 2 and tail[1] in {"me", "tenant", "validate", "refresh"}: return raise HTTPException(status_code=403, detail="Only WeKnora session introspection is allowed") if read_method and tail == ["me", "invitations", "pending-count"]: return if read_method and tail == ["tenants", "kv", "retrieval-config"]: return if tail[:1] == ["knowledge-bases"]: if len(tail) == 1: if method == "GET": return raise HTTPException(status_code=403, detail="Creating or deleting other WeKnora knowledge bases is not allowed") if tail[1] != ctx.weknora_id: raise HTTPException(status_code=403, detail="Only the current WeKnora knowledge base is allowed") if read_method or (write_method and ctx.can_write): return if tail[:1] == ["knowledgebase"]: if len(tail) < 2 or tail[1] != ctx.weknora_id: raise HTTPException(status_code=403, detail="Only the current WeKnora wiki is allowed") if read_method or (write_method and ctx.can_write): return if tail[:1] == ["knowledge"]: if len(tail) == 1: raise HTTPException(status_code=403, detail="Global WeKnora knowledge collection is not allowed") action = tail[1] if read_method and action == "batch": ids = [item.strip() for raw in request.query_params.getlist("ids") for item in raw.split(",") if item.strip()] if not ids: raise HTTPException(status_code=403, detail="Batch operation must include document ids") for document_id in ids: await _verify_proxy_document_belongs_to_kb(runtime, ctx, document_id, token) return if action in {"batch", "batch-delete", "batch-reparse", "folder", "move", "tags"}: if _request_kb_id(request, json_body) != ctx.weknora_id: raise HTTPException(status_code=403, detail="Batch operation must be scoped to the current knowledge base") return if action in {"search"}: raise HTTPException(status_code=403, detail="Global WeKnora search is disabled in embedded detail mode") document_id = tail[2] if action == "manual" and len(tail) >= 3 else action await _verify_proxy_document_belongs_to_kb(runtime, ctx, document_id, token) return if tail[:1] == ["chunks"]: if len(tail) < 2: raise HTTPException(status_code=403, detail="Chunk path is incomplete") if tail[1] == "by-id": if len(tail) < 3: raise HTTPException(status_code=403, detail="Chunk path is incomplete") await _verify_proxy_chunk_belongs_to_kb(runtime, ctx, tail[2], token) return await _verify_proxy_document_belongs_to_kb(runtime, ctx, tail[1], token) return if tail[:1] == ["knowledge-search"]: ids = json_body.get("knowledge_base_ids") if isinstance(json_body, dict) else None if ids == [ctx.weknora_id] or ids == (ctx.weknora_id,): return raise HTTPException(status_code=403, detail="Search must be scoped to the current WeKnora knowledge base") # A few read-only metadata endpoints are required by the original detail # page to render upload controls and current-tenant state. They do not # expose other knowledge-base contents. if ( read_method and tail and tail[0] in { "models", "system", "storage-backends", "vector-stores", "web-search-providers", "mcp-services", } ): return raise HTTPException(status_code=403, detail="This WeKnora endpoint is blocked in embedded detail mode") def _filter_knowledge_base_list_response(response: httpx.Response, ctx: WeKnoraEmbedContext) -> bytes | None: content_type = response.headers.get("content-type", "") if "application/json" not in content_type.lower(): return None try: payload = response.json() except ValueError: return None if not isinstance(payload, dict): return None allowed_weknora_ids = set(_embed_allowed_weknora_ids(ctx)) def keep_allowed(value: Any) -> Any: if isinstance(value, list): return [item for item in value if isinstance(item, dict) and str(item.get("id") or "") in allowed_weknora_ids] if isinstance(value, dict): result = dict(value) for key in ("items", "list", "knowledge_bases", "knowledge", "records", "data"): if isinstance(result.get(key), list): result[key] = keep_allowed(result[key]) result["total"] = len(result[key]) return result return value if isinstance(payload.get("data"), list): payload["data"] = keep_allowed(payload["data"]) payload["total"] = len(payload["data"]) elif isinstance(payload.get("data"), dict): payload["data"] = keep_allowed(payload["data"]) else: for key in ("items", "list", "knowledge_bases", "records"): if isinstance(payload.get(key), list): payload[key] = keep_allowed(payload[key]) payload["total"] = len(payload[key]) return json.dumps(payload, ensure_ascii=False, separators=(",", ":")).encode("utf-8") def _escape_script_json(value: Any) -> str: return json.dumps(value, ensure_ascii=False, separators=(",", ":")).replace(" str: raw = (value or "").strip() if not raw: return "" if "://" in raw: try: raw = urlsplit(raw).path except ValueError: return "" if not raw.startswith("/"): raw = f"/{raw}" raw = raw.rstrip("/") return "" if raw == "/" else raw def _request_public_prefix(request: Request) -> str: return _WEKNORA_PUBLIC_PREFIX def _with_public_prefix(prefix: str, path: str) -> str: if not path.startswith("/"): path = f"/{path}" return f"{prefix}{path}" if prefix else path def _without_public_prefix(prefix: str, path: str) -> str: if prefix and path.startswith(f"{prefix}/"): return path[len(prefix) :] if prefix and path == prefix: return "/" return path def _rewrite_weknora_html_public_paths( html: str, public_prefix: str, *, rewrite_relative: bool = True, ) -> str: if not public_prefix: return html for path in _WEKNORA_PUBLIC_ROOT_PATHS: prefixed = f"{public_prefix}{path}" html = re.sub( rf"(?P[\"'`(=:\s]){re.escape(path)}(?=[/?#\"'`),;\]\}}\s>]|$)", rf"\g{prefixed}", html, ) escaped_path = rf"\\/{re.escape(path.lstrip('/'))}" escaped_prefixed = f"\\/{prefixed.lstrip('/')}" html = re.sub( rf"(?]|$)", escaped_prefixed, html, ) if rewrite_relative: for path in _WEKNORA_PUBLIC_RELATIVE_ROOT_PATHS: prefixed = f"{public_prefix}/{path}" html = re.sub( rf"(?P[\"'`(=:\s])(?:\./)?{re.escape(path)}(?=/)", rf"\g{prefixed}", html, ) for path in _WEKNORA_PUBLIC_RELATIVE_FILE_PATHS: prefixed = f"{public_prefix}/{path}" html = re.sub( rf"(?P[\"'`(=:\s])(?:\./)?{re.escape(path)}(?=[?#\"'`),;\]\}}\s>]|$)", rf"\g{prefixed}", html, ) return html def _rewrite_weknora_script_public_paths(text: str, public_prefix: str) -> str: if not public_prefix: return text for path in _WEKNORA_PUBLIC_ROOT_PATHS: prefixed = f"{public_prefix}{path}" text = re.sub( rf"(?P[\"'`]){re.escape(path)}(?=[/?#\"'`),;\]\}}\s]|$)", rf"\g{prefixed}", text, ) return _rewrite_weknora_vite_dependency_map(text, public_prefix) def _should_rewrite_weknora_static_response(content_type: str, upstream_path: str) -> bool: lower_type = content_type.lower() if any( marker in lower_type for marker in ( "javascript", "text/css", "application/json", "manifest", "image/svg+xml", ) ): return True lower_path = upstream_path.lower() return lower_path.endswith((".js", ".css", ".json", ".webmanifest", ".svg")) def _rewrite_weknora_vite_dependency_map(text: str, public_prefix: str) -> str: if not public_prefix or "__vite__mapDeps" not in text: return text marker = "m.f||(m.f=[" start = text.find(marker) if start < 0: return text end = text.find("]))", start + len(marker)) if end < 0: return text prefix = public_prefix.strip("/") dependency_map = text[start:end] dependency_map = re.sub( r'(?P["\'])assets/', rf"\g{prefix}/assets/", dependency_map, ) return f"{text[:start]}{dependency_map}{text[end:]}" def _weknora_static_fallback_upstream_path(request: Request) -> str | None: path = _without_public_prefix( _request_public_prefix(request), str(request.scope.get("path") or request.url.path), ) if not path.startswith("/"): path = f"/{path}" lower_path = path.lower() if lower_path.startswith("/api/"): return None if lower_path.startswith(_WEKNORA_STATIC_ROOT_PREFIXES): return path for root_prefix in _WEKNORA_STATIC_ROOT_PREFIXES: index = lower_path.find(root_prefix) if index > 0: return path[index:] for asset_file in _WEKNORA_ROOT_ASSET_FILES: if lower_path.endswith(asset_file): return asset_file if lower_path.endswith(_WEKNORA_STATIC_FILE_SUFFIXES): return path return None def _rewrite_weknora_location_header(location: str, public_prefix: str, target_base_url: str) -> str: if not location or not public_prefix: return location try: current = urlsplit(location) upstream = urlsplit(target_base_url) except ValueError: return location def proxied_path(path: str) -> str: upstream_base_path = upstream.path.rstrip("/") if upstream_base_path and (path == upstream_base_path or path.startswith(f"{upstream_base_path}/")): path = path[len(upstream_base_path) :] or "/" return _without_public_prefix(public_prefix, path) suffix = "" if current.query: suffix += f"?{current.query}" if current.fragment: suffix += f"#{current.fragment}" if not current.scheme and location.startswith("/"): path = proxied_path(current.path or "/") return f"{public_prefix}{path}{suffix}" if current.scheme and current.netloc and current.netloc == upstream.netloc: path = proxied_path(current.path or "/") return f"{public_prefix}{path}{suffix}" return location def _weknora_embed_injection( *, allowed_path: str, tab: str, session: dict[str, Any], ctx: WeKnoraEmbedContext, public_prefix: str = "", ) -> str: tenant_id = _tenant_id_from_session(session) user = session.get("user") if isinstance(session.get("user"), dict) else {} tenant = session.get("tenant") if isinstance(session.get("tenant"), dict) else {} payload = { "allowedPath": allowed_path, "tab": tab, "canWrite": ctx.can_write, "token": session.get("token") or "", "refreshToken": session.get("refresh_token") or "", "tenantId": tenant_id, "user": user, "tenant": tenant, "weknoraId": ctx.weknora_id, "allowedWeKnoraIds": list(_embed_allowed_weknora_ids(ctx)), "publicPrefix": public_prefix, "apiPrefix": _with_public_prefix(public_prefix, _WEKNORA_EMBED_API_PREFIX), } data = _escape_script_json(payload) return f""" """ def _inject_weknora_embed_html( html: bytes, *, allowed_path: str, tab: str, session: dict[str, Any], ctx: WeKnoraEmbedContext, public_prefix: str = "", ) -> bytes: text = html.decode("utf-8", errors="replace") text = _rewrite_weknora_html_public_paths(text, public_prefix) injection = _weknora_embed_injection( allowed_path=allowed_path, tab=tab, session=session, ctx=ctx, public_prefix=public_prefix, ) if "" in text: text = text.replace("", f"{injection}", 1) else: text = f"{injection}{text}" return text.encode("utf-8") async def _proxy_weknora_request( request: Request, *, target_base_url: str, upstream_path: str, ctx: WeKnoraEmbedContext, is_api: bool, tab: str = "wiki", ) -> Response: runtime = _runtime_for_iframe_proxy(request) session = await _get_weknora_admin_session(runtime) token = str(session.get("token") or "") body = await request.body() if is_api: await _authorize_weknora_api_proxy(request, ctx, runtime, upstream_path, body, token) headers = _forward_request_headers(request) if is_api: headers["Authorization"] = f"Bearer {token}" tenant_id = _tenant_id_from_session(session) if tenant_id: headers["X-Tenant-ID"] = tenant_id url = f"{target_base_url.rstrip('/')}{_encode_upstream_path(upstream_path)}" follow_binary_redirects = is_api and _should_follow_weknora_binary_redirect(upstream_path) try: async with httpx.AsyncClient( timeout=httpx.Timeout(60.0), follow_redirects=follow_binary_redirects, trust_env=False, ) as client: response = await client.request( request.method, url, params=request.query_params.multi_items(), headers=headers, content=body, ) if is_api and response.status_code in {401, 403}: session = await _get_weknora_admin_session(runtime, force_refresh=True) headers["Authorization"] = f"Bearer {session.get('token')}" tenant_id = _tenant_id_from_session(session) if tenant_id: headers["X-Tenant-ID"] = tenant_id async with httpx.AsyncClient( timeout=httpx.Timeout(60.0), follow_redirects=follow_binary_redirects, trust_env=False, ) as client: response = await client.request( request.method, url, params=request.query_params.multi_items(), headers=headers, content=body, ) except httpx.TimeoutException as exc: raise HTTPException(status_code=504, detail="WeKnora proxy request timed out") from exc except httpx.HTTPError as exc: logger.warning("WeKnora proxy request failed: url=%s error=%r", url, exc) raise HTTPException(status_code=502, detail="WeKnora proxy request failed") from exc content = response.content content_type = response.headers.get("content-type", "") public_prefix = _request_public_prefix(request) if is_api and request.method.upper() == "GET" and upstream_path.rstrip("/") == "/api/v1/knowledge-bases": content = _filter_knowledge_base_list_response(response, ctx) or content if is_api and ctx.is_conversation_deposit and "application/json" in content_type.lower(): try: payload = json.loads(content) except (TypeError, ValueError, json.JSONDecodeError): pass else: content = json.dumps( _sanitize_conversation_deposit_value(payload), ensure_ascii=False, separators=(",", ":"), ).encode("utf-8") if not is_api and "text/html" in content_type.lower(): allowed_path = _with_public_prefix( public_prefix, f"/platform/knowledge-bases/{quote(ctx.weknora_id, safe='')}", ) content = _inject_weknora_embed_html( content, allowed_path=allowed_path, tab=tab, session=session, ctx=ctx, public_prefix=public_prefix, ) content_type = "text/html; charset=utf-8" elif not is_api and _should_rewrite_weknora_static_response(content_type, upstream_path): text = content.decode("utf-8", errors="replace") lower_type = content_type.lower() if "javascript" in lower_type or upstream_path.lower().endswith((".js", ".mjs")): text = _rewrite_weknora_script_public_paths(text, public_prefix) else: text = _rewrite_weknora_html_public_paths( text, public_prefix, rewrite_relative=False, ) content = text.encode("utf-8") headers_out = _safe_response_headers(response.headers) if "location" in headers_out: headers_out["location"] = _rewrite_weknora_location_header( headers_out["location"], public_prefix, target_base_url, ) if "Location" in headers_out: headers_out["Location"] = _rewrite_weknora_location_header( headers_out["Location"], public_prefix, target_base_url, ) if not is_api: headers_out.pop("etag", None) headers_out.pop("ETag", None) headers_out["Cache-Control"] = "no-store" if is_api or "text/html" in content_type.lower() else "private, max-age=3600" return Response( content=content, status_code=response.status_code, media_type=content_type.split(";", 1)[0] if content_type else None, headers=headers_out, ) def _is_conversation_deposit_mapping(row: dict[str, Any]) -> bool: return str(row.get("owner_user_id") or "") == CONVERSATION_DEPOSIT_OWNER_USER_ID and str(row.get("name") or "").strip() == CONVERSATION_DEPOSIT_KB_NAME _LEGACY_BRAND_PATTERN = re.compile(r"deer[\s_-]*flow", re.IGNORECASE) def _sanitize_conversation_deposit_value(value: Any) -> Any: if isinstance(value, str): return _LEGACY_BRAND_PATTERN.sub("cmzs", value) if isinstance(value, dict): return {key: _sanitize_conversation_deposit_value(item) for key, item in value.items()} if isinstance(value, list): return [_sanitize_conversation_deposit_value(item) for item in value] if isinstance(value, tuple): return tuple(_sanitize_conversation_deposit_value(item) for item in value) return value async def _sync_conversation_deposit_view( store: LlmWikiStore, client: WeKnoraClient, row: dict[str, Any], remote: dict[str, Any] | None, ) -> tuple[dict[str, Any], dict[str, Any] | None]: if not _is_conversation_deposit_mapping(row): return row, remote current_remote = remote if remote is not None and str(remote.get("description") or "") != CONVERSATION_DEPOSIT_KB_DESCRIPTION: try: updated_remote = await client.update_knowledge_base( str(row["weknora_id"]), {"description": CONVERSATION_DEPOSIT_KB_DESCRIPTION}, ) current_remote = { **remote, **updated_remote, "description": CONVERSATION_DEPOSIT_KB_DESCRIPTION, } except WeKnoraError: logger.warning("Failed to synchronize the cmzs conversation-deposit description", exc_info=True) current_row = row if str(row.get("description") or "") != CONVERSATION_DEPOSIT_KB_DESCRIPTION: updated_row = await store.update_mapping( str(row["id"]), description=CONVERSATION_DEPOSIT_KB_DESCRIPTION, ) current_row = updated_row or row return current_row, current_remote def _mapping_view(row: dict[str, Any], remote: dict[str, Any] | None, *, actor_user_id: str, is_admin: bool) -> dict[str, Any]: source = remote or {} is_conversation_deposit = _is_conversation_deposit_mapping(row) is_owner = row.get("owner_user_id") == actor_user_id description = CONVERSATION_DEPOSIT_KB_DESCRIPTION if is_conversation_deposit else source.get("description") if source.get("description") is not None else row.get("description") or "" return { "id": row["id"], "remote_status": "available" if remote is not None else "missing", "name": source.get("name") or row.get("name") or "", "description": description, "type": source.get("type") or row.get("kb_type") or "document", "publication_status": "published" if is_conversation_deposit else row.get("publication_status") or "private", "wiki_index_enabled": row.get("wiki_index_enabled", True), "external_search_enabled": False if is_conversation_deposit else row.get("external_search_enabled", False), "external_search_updated_at": row.get("external_search_updated_at"), "owner_user_id": row.get("owner_user_id") if is_admin else None, "is_owner": is_owner, "can_write": not is_conversation_deposit and (is_admin or is_owner), "is_conversation_deposit": is_conversation_deposit, "knowledge_count": int(source.get("knowledge_count") or 0), "chunk_count": int(source.get("chunk_count") or 0), "processing_count": int(source.get("processing_count") or 0), "share_count": int(source.get("share_count") or 0), "is_processing": bool(source.get("is_processing") or int(source.get("processing_count") or 0) > 0), "is_temporary": bool(source.get("is_temporary", False)), "embedding_model_id": source.get("embedding_model_id"), "vector_store_id": source.get("vector_store_id"), "vector_store_name": source.get("vector_store_name"), "vector_store_source": source.get("vector_store_source"), "vector_store_engine_type": source.get("vector_store_engine_type"), "vector_store_status": source.get("vector_store_status"), "chunking_config": source.get("chunking_config") if isinstance(source.get("chunking_config"), dict) else {}, "image_processing_config": (source.get("image_processing_config") if isinstance(source.get("image_processing_config"), dict) else {}), "indexing_strategy": source.get("indexing_strategy") if isinstance(source.get("indexing_strategy"), dict) else {}, "capabilities": source.get("capabilities") if isinstance(source.get("capabilities"), dict) else {}, "extract_config": source.get("extract_config") if isinstance(source.get("extract_config"), dict) else {}, "faq_config": source.get("faq_config") if isinstance(source.get("faq_config"), dict) else {}, "wiki_config": source.get("wiki_config") if isinstance(source.get("wiki_config"), dict) else {}, "question_generation_config": (source.get("question_generation_config") if isinstance(source.get("question_generation_config"), dict) else {}), "auto_tag_config": source.get("auto_tag_config") if isinstance(source.get("auto_tag_config"), dict) else {}, "vlm_config": source.get("vlm_config") if isinstance(source.get("vlm_config"), dict) else {}, "asr_config": source.get("asr_config") if isinstance(source.get("asr_config"), dict) else {}, "summary_model_id": source.get("summary_model_id"), "storage_provider": (source.get("storage_provider_config", {}).get("provider") if isinstance(source.get("storage_provider_config"), dict) else None), "created_at": row.get("created_at"), "updated_at": row.get("updated_at"), "published_at": row.get("published_at"), "remote_created_at": source.get("created_at"), "remote_updated_at": source.get("updated_at"), } async def _authorized_mapping( request: Request, mapping_id: str, *, write: bool, ) -> tuple[dict[str, Any], str, bool]: user_id, is_admin = await _actor(request) row = await _store(request).get_authorized(mapping_id, user_id, write=write, is_admin=is_admin) if row is None: raise HTTPException(status_code=404, detail="Knowledge base not found") if write and _is_conversation_deposit_mapping(row): raise HTTPException(status_code=403, detail="Conversation deposit is managed by cmzs") return row, user_id, is_admin async def _public_mapping(request: Request, mapping_id: str) -> dict[str, Any]: """Resolve a world-readable knowledge base without inheriting actor/admin state.""" row = await _store(request).get_authorized( mapping_id, "__public_wiki__", write=False, is_admin=False, ) if row is None or row.get("publication_status") != "published" or _is_conversation_deposit_mapping(row): raise HTTPException(status_code=404, detail="Published knowledge base not found") return row def _validate_protected_wiki_file_path(file_path: str) -> str: value = file_path.strip() if not value or "\x00" in value or not _PROTECTED_WIKI_FILE_PATH.match(value): raise HTTPException(status_code=422, detail="Unsupported Wiki image path") return value async def _wiki_file_response(row: dict[str, Any], file_path: str, *, public: bool) -> Response: safe_file_path = _validate_protected_wiki_file_path(file_path) try: content, media_type = await _client_or_503().get_knowledge_base_file( str(row["weknora_id"]), safe_file_path, ) except WeKnoraError as exc: _raise_upstream(exc) filename = unquote(safe_file_path.rstrip("/").rsplit("/", 1)[-1]).replace("\r", "").replace("\n", "") or "wiki-file" if "." not in filename: filename += mimetypes.guess_extension(str(media_type).split(";", 1)[0].strip()) or "" return Response( content=content, media_type=media_type, headers={ "Cache-Control": "public, max-age=300" if public else "private, max-age=3600", "Content-Disposition": f"inline; filename*=UTF-8''{quote(filename, safe='')}", "X-Content-Type-Options": "nosniff", }, ) async def _verify_document(client: WeKnoraClient, row: dict[str, Any], document_id: str) -> dict[str, Any]: document = await client.get_document(document_id) if str(document.get("knowledge_base_id") or "") != str(row["weknora_id"]): raise HTTPException(status_code=404, detail="Document not found") return document def _document_view(document: dict[str, Any]) -> dict[str, Any]: """Expose useful document state without leaking WeKnora storage paths or hashes.""" fields = ( "id", "type", "title", "description", "source", "channel", "parse_status", "pending_subtasks_count", "summary_status", "enable_status", "embedding_model_id", "file_name", "file_type", "file_size", "storage_size", "created_at", "updated_at", "processed_at", "error_message", ) return {key: document.get(key) for key in fields} def _related_chunk_ids(value: Any) -> list[str]: if not isinstance(value, list): return [] result: list[str] = [] for item in value: candidate = item.get("id") if isinstance(item, dict) else item if candidate and str(candidate) not in result: result.append(str(candidate)) return result def _chunk_view(chunk: dict[str, Any]) -> dict[str, Any]: return { "id": str(chunk.get("id") or ""), "content": str(chunk.get("content") or ""), "chunk_index": int(chunk.get("chunk_index") or 0), "is_enabled": bool(chunk.get("is_enabled")), "status": chunk.get("status"), "start_at": int(chunk.get("start_at") or 0), "end_at": int(chunk.get("end_at") or 0), "chunk_type": str(chunk.get("chunk_type") or "text"), "parent_chunk_id": str(chunk.get("parent_chunk_id") or ""), "relation_chunk_ids": _related_chunk_ids(chunk.get("relation_chunks")), "indirect_relation_chunk_ids": _related_chunk_ids(chunk.get("indirect_relation_chunks")), "metadata": chunk.get("metadata") if isinstance(chunk.get("metadata"), dict) else {}, "created_at": chunk.get("created_at"), "updated_at": chunk.get("updated_at"), } def _string_tokens(value: Any) -> list[str]: tokens: list[str] = [] if isinstance(value, str): raw = value.strip() if raw: tokens.append(raw) for separator in ("#", ":", "/", "\\"): if separator in raw: tokens.extend(part.strip() for part in raw.split(separator) if part.strip()) elif isinstance(value, dict): for key in ("id", "chunk_id", "knowledge_id", "document_id", "source_id", "ref_id", "slug"): tokens.extend(_string_tokens(value.get(key))) return list(dict.fromkeys(tokens)) def _source_ref_tokens(value: Any) -> set[str]: refs: set[str] = set() if isinstance(value, list): for item in value: refs.update(_string_tokens(item)) else: refs.update(_string_tokens(value)) return refs def _candidate_chunk_ref_tokens(chunk: dict[str, Any], chunk_id: str) -> set[str]: refs = {chunk_id, str(chunk.get("id") or "")} metadata = chunk.get("metadata") if isinstance(chunk.get("metadata"), dict) else {} for source in (chunk, metadata): for key in ( "knowledge_id", "document_id", "source_id", "source_document_id", "knowledge_uuid", "file_id", ): refs.update(_string_tokens(source.get(key))) return {item for item in refs if item} def _candidate_wiki_slugs(chunk: dict[str, Any]) -> list[str]: metadata = chunk.get("metadata") if isinstance(chunk.get("metadata"), dict) else {} candidates: list[str] = [] for source in (chunk, metadata): for key in ("wiki_page_slug", "page_slug", "wiki_slug", "llmwiki_slug", "slug"): candidates.extend(_string_tokens(source.get(key))) nested_wiki = metadata.get("wiki") if isinstance(nested_wiki, dict): for key in ("slug", "page_slug"): candidates.extend(_string_tokens(nested_wiki.get(key))) return list(dict.fromkeys(candidates)) def _wiki_page_view(page: dict[str, Any]) -> dict[str, Any] | None: slug = str(page.get("slug") or "").strip() title = str(page.get("title") or slug or "").strip() content = str(page.get("content") or "").strip() summary = str(page.get("summary") or "").strip() if not slug and not title and not content and not summary: return None return { "page_metadata": page.get("page_metadata") if isinstance(page.get("page_metadata"), dict) else {}, "id": str(page.get("id") or ""), "slug": slug, "title": title, "page_type": str(page.get("page_type") or "page"), "status": str(page.get("status") or ""), "summary": summary, "content": content, "aliases": list(page.get("aliases") or []) if isinstance(page.get("aliases"), list) else [], "parent_slug": str(page.get("parent_slug") or page.get("parentSlug") or ""), "category_path": list(page.get("category_path") or []) if isinstance(page.get("category_path"), list) else [], "folder_id": str(page.get("folder_id") or page.get("folderId") or ""), "wiki_path": str(page.get("wiki_path") or page.get("wikiPath") or ""), "depth": int(page.get("depth") or 0), "version": int(page.get("version") or 0), "created_at": str(page.get("created_at") or ""), "updated_at": str(page.get("updated_at") or ""), "source_refs": list(page.get("source_refs") or []) if isinstance(page.get("source_refs"), list) else [], "in_links": list(page.get("in_links") or []) if isinstance(page.get("in_links"), list) else [], "out_links": list(page.get("out_links") or []) if isinstance(page.get("out_links"), list) else [], } async def _all_wiki_pages(client: WeKnoraClient, remote_id: str) -> list[dict[str, Any]]: """Read every Wiki page and its processed Markdown, never raw documents.""" summaries: dict[str, dict[str, Any]] = {} page_number = 1 while True: listing = await client.list_wiki_pages(remote_id, page=page_number, page_size=500) rows = [item for item in listing.get("pages", []) if isinstance(item, dict)] for item in rows: slug = str(item.get("slug") or item.get("wiki_slug") or item.get("path") or "").strip("/") if slug: summaries[slug] = item total_pages = max(1, int(listing.get("total_pages") or 1)) if page_number >= total_pages or not rows: break page_number += 1 pages: list[dict[str, Any]] = [] entries = list(summaries.items()) for start in range(0, len(entries), 20): batch = entries[start : start + 20] details = await asyncio.gather(*(client.get_wiki_page(remote_id, slug) for slug, _ in batch)) pages.extend({**summary, **detail, "slug": slug} for (slug, summary), detail in zip(batch, details, strict=True)) return pages def _wiki_directory_parts(page: dict[str, Any], by_slug: dict[str, dict[str, Any]]) -> list[str]: parent_slug = str(page.get("parent_slug") or "").strip("/") if parent_slug: ancestors: list[str] = [] visited: set[str] = set() current = by_slug.get(parent_slug) while current is not None: slug = str(current.get("slug") or "").strip("/") if not slug or slug in visited: break visited.add(slug) ancestors.insert(0, str(current.get("title") or slug.rsplit("/", 1)[-1])) current_parent = str(current.get("parent_slug") or "").strip("/") current = by_slug.get(current_parent) if current_parent else None if ancestors: return ancestors category_path = [str(item).strip() for item in page.get("category_path") or [] if str(item).strip()] if category_path: return category_path wiki_path = [item.strip() for item in str(page.get("wiki_path") or "").split("/") if item.strip()] if len(wiki_path) > 2: return wiki_path[1:-1] slug_parts = [item for item in str(page.get("slug") or "").split("/") if item] return slug_parts[:-1] def _safe_excel_text(value: Any) -> str: text = str(value or "") return f"'{text}" if text.startswith(("=", "+", "-", "@")) else text def _build_wiki_excel(knowledge_base_name: str, raw_pages: list[dict[str, Any]]) -> BytesIO: from openpyxl import Workbook from openpyxl.styles import Alignment, Font, PatternFill pages = [view for raw in raw_pages if (view := _wiki_page_view(raw)) is not None] by_slug = {page["slug"]: page for page in pages if page["slug"]} directories = [_wiki_directory_parts(page, by_slug) for page in pages] max_depth = max((len(parts) for parts in directories), default=0) max_content_parts = max((max(1, (len(str(page.get("content") or "")) + 31_999) // 32_000) for page in pages), default=1) columns = [ "知识库", "标题", "Slug", "目录完整路径", "所在目录", "父级目录完整路径", *(f"第{level}级目录" for level in range(1, max_depth + 1)), "摘要", *("正文" if index == 1 else f"正文(续{index - 1})" for index in range(1, max_content_parts + 1)), "页面类型", "状态", "别名", "版本", "创建时间", "更新时间", ] workbook = Workbook() sheet = workbook.active sheet.title = "Wiki 全量导出" sheet.append(columns) for cell in sheet[1]: cell.font = Font(bold=True, color="FFFFFF") cell.fill = PatternFill("solid", fgColor="1F4E78") cell.alignment = Alignment(horizontal="center", vertical="center") for page, directory in zip(pages, directories, strict=True): content = str(page.get("content") or "") content_parts = [content[index : index + 32_000] for index in range(0, len(content), 32_000)] or [""] content_parts.extend([""] * (max_content_parts - len(content_parts))) row = [ knowledge_base_name, page.get("title"), page.get("slug"), " / ".join(directory), directory[-1] if directory else "根目录", " / ".join(directory[:-1]), *directory, *([""] * (max_depth - len(directory))), page.get("summary"), *content_parts, page.get("page_type"), page.get("status"), "、".join(str(item) for item in page.get("aliases") or []), page.get("version"), page.get("created_at"), page.get("updated_at"), ] sheet.append([_safe_excel_text(value) for value in row]) sheet.freeze_panes = "A2" sheet.auto_filter.ref = sheet.dimensions for column in sheet.columns: letter = column[0].column_letter heading = str(column[0].value or "") sheet.column_dimensions[letter].width = 60 if heading.startswith("正文") else 32 if heading in {"摘要", "目录完整路径", "父级目录完整路径"} else 20 for cell in column[1:]: cell.alignment = Alignment(vertical="top", wrap_text=True) output = BytesIO() workbook.save(output) output.seek(0) return output def _weknora_embed_target_path(weknora_id: str, tab: str) -> str: target_query = urlencode({"tab": tab}) if tab != "documents" else "" target = f"/platform/knowledge-bases/{quote(weknora_id, safe='')}" if target_query: target = f"{target}?{target_query}" return target def _set_weknora_embed_cookie(request: Request, response: Response, payload: dict[str, Any]) -> None: secure = is_secure_request(request) signed = _sign_embed_payload(payload) for path in ("/", _WEKNORA_PUBLIC_PREFIX): response.set_cookie( _WEKNORA_EMBED_COOKIE, signed, max_age=_WEKNORA_EMBED_TTL_SECONDS, httponly=True, secure=secure, samesite="none" if secure else "lax", path=path, ) response.headers["Cache-Control"] = "no-store" async def _find_wiki_page_for_chunk( client: WeKnoraClient, *, knowledge_base_id: str, chunk_id: str, chunks: list[dict[str, Any]], ) -> dict[str, Any] | None: focus = next((chunk for chunk in chunks if str(chunk.get("id") or "") == chunk_id), None) if focus is None: return None for slug in _candidate_wiki_slugs(focus): try: page = await client.get_wiki_page(knowledge_base_id, slug) except WeKnoraError as exc: if exc.status_code != 404: logger.debug("Failed to read WeKnora wiki page slug=%s", slug, exc_info=True) continue view = _wiki_page_view(page) if view and (view.get("content") or view.get("summary")): return view return None async def _prepare_weknora_embed_session( request: Request, mapping_id: str, ) -> tuple[str, dict[str, Any]]: _runtime_for_iframe_proxy(request) row, user_id, is_admin = await _authorized_mapping(request, mapping_id, write=False) visible_rows = await _store(request).list_visible( user_id, scope="all", # The iframe must mirror the knowledge-base picker, not grant an # administrator an unlisted user's private knowledge bases. is_admin=False, ) writable = None if not _is_conversation_deposit_mapping(row): writable = await _store(request).get_authorized(mapping_id, user_id, write=True, is_admin=is_admin) weknora_id = str(row["weknora_id"]) ctx = WeKnoraEmbedContext( mapping_id=mapping_id, weknora_id=weknora_id, user_id=user_id, is_admin=is_admin, can_write=writable is not None, is_conversation_deposit=_is_conversation_deposit_mapping(row), allowed_weknora_ids=tuple(dict.fromkeys(value for value in (weknora_id, *(str(item.get("weknora_id") or "") for item in visible_rows)) if value)), exp=int(time.time()) + _WEKNORA_EMBED_TTL_SECONDS, ) return weknora_id, _embed_context_payload(ctx) @router.post("/api/llmwiki/knowledge-bases/{mapping_id}/weknora/session") async def create_weknora_detail_session( request: Request, mapping_id: str, tab: Literal["documents", "wiki", "graph"] = Query(default="wiki"), ) -> JSONResponse: """Create a signed iframe session using the normal DeerFlow API auth path.""" weknora_id, payload = await _prepare_weknora_embed_session(request, mapping_id) target = _with_public_prefix( _request_public_prefix(request), _weknora_embed_target_path(weknora_id, tab), ) response = JSONResponse( { "frame_path": target, "expires_in": _WEKNORA_EMBED_TTL_SECONDS, } ) _set_weknora_embed_cookie(request, response, payload) return response @router.get("/api/llmwiki/knowledge-bases/{mapping_id}/weknora/frame") async def open_weknora_detail_frame( request: Request, mapping_id: str, tab: Literal["documents", "wiki", "graph"] = Query(default="wiki"), ) -> RedirectResponse: """Start a short-lived, DeerFlow-authorized WeKnora detail iframe session.""" weknora_id, payload = await _prepare_weknora_embed_session(request, mapping_id) target = _with_public_prefix( _request_public_prefix(request), _weknora_embed_target_path(weknora_id, tab), ) response = RedirectResponse(target, status_code=302) _set_weknora_embed_cookie(request, response, payload) return response @proxy_router.api_route("/platform", methods=["GET", "HEAD"]) @proxy_router.api_route("/platform/{proxied_path:path}", methods=["GET", "HEAD"]) async def proxy_weknora_platform_page(request: Request, proxied_path: str = "") -> Response: ctx = await _read_embed_context(request) requested = f"/platform/{proxied_path}".rstrip("/") allowed_path = f"/platform/knowledge-bases/{quote(ctx.weknora_id, safe='')}" if requested != allowed_path: requested_weknora_id = _weknora_id_from_platform_path(requested) if requested_weknora_id is not None: selected_ctx = await _switch_embed_context(request, ctx, requested_weknora_id) response = await _proxy_weknora_request( request, target_base_url=_runtime_for_iframe_proxy(request).web_base_url, upstream_path=f"/platform/knowledge-bases/{quote(requested_weknora_id, safe='')}", ctx=selected_ctx, is_api=False, tab=request.query_params.get("tab") or "wiki", ) _set_weknora_embed_cookie(request, response, _embed_context_payload(selected_ctx)) return response static_upstream_path = _weknora_static_fallback_upstream_path(request) if static_upstream_path is not None: runtime = _runtime_for_iframe_proxy(request) return await _proxy_weknora_request( request, target_base_url=runtime.web_base_url, upstream_path=static_upstream_path, ctx=ctx, is_api=False, ) target_query = urlencode({"tab": request.query_params.get("tab") or "wiki"}) public_target = _with_public_prefix(_request_public_prefix(request), allowed_path) return RedirectResponse(f"{public_target}?{target_query}", status_code=302) runtime = _runtime_for_iframe_proxy(request) return await _proxy_weknora_request( request, target_base_url=runtime.web_base_url, upstream_path=allowed_path, ctx=ctx, is_api=False, tab=request.query_params.get("tab") or "wiki", ) @proxy_router.api_route("/assets/{proxied_path:path}", methods=["GET", "HEAD"]) async def proxy_weknora_assets(request: Request, proxied_path: str) -> Response: ctx = await _read_embed_context(request) runtime = _runtime_for_iframe_proxy(request) return await _proxy_weknora_request( request, target_base_url=runtime.web_base_url, upstream_path=f"/assets/{proxied_path}", ctx=ctx, is_api=False, ) @proxy_router.api_route("/locales/{proxied_path:path}", methods=["GET", "HEAD"]) async def proxy_weknora_locales(request: Request, proxied_path: str) -> Response: ctx = await _read_embed_context(request) runtime = _runtime_for_iframe_proxy(request) return await _proxy_weknora_request( request, target_base_url=runtime.web_base_url, upstream_path=f"/locales/{proxied_path}", ctx=ctx, is_api=False, ) @proxy_router.api_route("/favicon.ico", methods=["GET", "HEAD"]) @proxy_router.api_route("/favicon.svg", methods=["GET", "HEAD"]) @proxy_router.api_route("/config.js", methods=["GET", "HEAD"]) @proxy_router.api_route("/logo.svg", methods=["GET", "HEAD"]) @proxy_router.api_route("/manifest.webmanifest", methods=["GET", "HEAD"]) async def proxy_weknora_root_asset(request: Request) -> Response: ctx = await _read_embed_context(request) runtime = _runtime_for_iframe_proxy(request) return await _proxy_weknora_request( request, target_base_url=runtime.web_base_url, upstream_path=_without_public_prefix( _request_public_prefix(request), str(request.scope.get("path") or request.url.path), ), ctx=ctx, is_api=False, ) @proxy_router.api_route("/tdesign-icons/{proxied_path:path}", methods=["GET", "HEAD"]) async def proxy_weknora_tdesign_icons(request: Request, proxied_path: str) -> Response: ctx = await _read_embed_context(request) runtime = _runtime_for_iframe_proxy(request) return await _proxy_weknora_request( request, target_base_url=runtime.web_base_url, upstream_path=f"/tdesign-icons/{proxied_path}", ctx=ctx, is_api=False, ) @proxy_router.api_route( "/api/v1/{proxied_path:path}", methods=["GET", "HEAD", "POST", "PUT", "PATCH", "DELETE", "OPTIONS"], ) async def proxy_weknora_api(request: Request, proxied_path: str) -> Response: if request.method.upper() == "OPTIONS": return Response(status_code=204) ctx = await _read_embed_context(request) runtime = _runtime_for_iframe_proxy(request) return await _proxy_weknora_request( request, target_base_url=runtime.api_base_url, upstream_path=f"/api/v1/{proxied_path}", ctx=ctx, is_api=True, ) @router.api_route( "/api/llmwiki/weknora-embed/api/v1/{proxied_path:path}", methods=["GET", "HEAD", "POST", "PUT", "PATCH", "DELETE", "OPTIONS"], ) @router.api_route( "/deerflow/api/llmwiki/weknora-embed/api/v1/{proxied_path:path}", methods=["GET", "HEAD", "POST", "PUT", "PATCH", "DELETE", "OPTIONS"], ) async def proxy_weknora_embed_api(request: Request, proxied_path: str) -> Response: """Proxy iframe API calls under a collision-free cmzs namespace.""" if request.method.upper() == "OPTIONS": return Response(status_code=204) ctx = await _read_embed_context(request) runtime = _runtime_for_iframe_proxy(request) return await _proxy_weknora_request( request, target_base_url=runtime.api_base_url, upstream_path=f"/api/v1/{proxied_path}", ctx=ctx, is_api=True, ) @proxy_router.api_route("/{proxied_path:path}", methods=["GET", "HEAD"]) async def proxy_weknora_static_fallback(request: Request, proxied_path: str) -> Response: upstream_path = _weknora_static_fallback_upstream_path(request) if upstream_path is None: raise HTTPException(status_code=404, detail="WeKnora proxy path not found") ctx = await _read_embed_context(request) runtime = _runtime_for_iframe_proxy(request) return await _proxy_weknora_request( request, target_base_url=runtime.web_base_url, upstream_path=upstream_path, ctx=ctx, is_api=False, ) @router.get("/api/llmwiki/runtime") async def get_llmwiki_runtime() -> dict[str, Any]: return _runtime_or_legacy() @router.get("/api/system-settings/llmwiki") async def get_llmwiki_settings(request: Request) -> dict[str, Any]: await _require_admin(request) app_config = request.app.state.config runtime = get_resolved_llmwiki_runtime(app_config) return { **_runtime_or_legacy(), "file_api_base_url": app_config.llmwiki.weknora.api_base_url, "file_web_base_url": app_config.llmwiki.weknora.web_base_url, "effective_provider": runtime.provider, } @router.put("/api/system-settings/llmwiki") async def update_llmwiki_settings(request: Request, body: LlmWikiSettingsUpdate) -> dict[str, Any]: await _require_admin(request) try: override = LlmWikiRuntimeOverride( enabled=body.override_enabled, api_base_url=body.api_base_url, web_base_url=body.web_base_url, ) except ValueError as exc: raise HTTPException(status_code=422, detail=str(exc)) from None current = load_system_settings() save_system_settings(current.model_copy(update={"llmwiki": override})) return _runtime_or_legacy() @router.delete("/api/system-settings/llmwiki/override") async def reset_llmwiki_settings(request: Request) -> dict[str, Any]: await _require_admin(request) current = load_system_settings() save_system_settings(current.model_copy(update={"llmwiki": LlmWikiRuntimeOverride()})) return _runtime_or_legacy() @router.post("/api/system-settings/llmwiki/test") async def test_llmwiki_connection(request: Request, body: LlmWikiConnectionTest) -> dict[str, Any]: await _require_admin(request) runtime = get_resolved_llmwiki_runtime(request.app.state.config) if body.api_base_url is not None: try: candidate = normalize_service_url(body.api_base_url) except ValueError as exc: raise HTTPException(status_code=422, detail=str(exc)) from None runtime = replace(runtime, provider="weknora" if candidate else "legacy", api_base_url=candidate) if not runtime.weknora_enabled: raise HTTPException(status_code=422, detail="Enter a WeKnora API address first") try: health = await build_weknora_client(runtime).health() except RuntimeError as exc: raise HTTPException(status_code=503, detail=str(exc)) from None except WeKnoraError as exc: _raise_upstream(exc) return {"success": True, "provider": "weknora", "health": health} _PUBLIC_WIKI_SUMMARY_CONCURRENCY = 8 _PUBLIC_WIKI_PAGE_SIZE = 500 _PUBLIC_WIKI_PAGE_LIMIT = 100 _GENERATED_PUBLIC_WIKI_INDEX_SLUG = "__index__" def _is_public_wiki_index_page(page: dict[str, Any]) -> bool: page_type = str(page.get("page_type") or "").strip().lower() title = str(page.get("title") or "").strip().lower() slug_name = str(page.get("slug") or "").strip("/").rsplit("/", 1)[-1].lower() return page_type == "index" or title in {"索引", "index"} or slug_name in {"index", "home", "readme"} def _public_wiki_default_page(pages: list[dict[str, Any]]) -> dict[str, Any] | None: return ( next((page for page in pages if _is_public_wiki_index_page(page)), None) or next((page for page in pages if str(page.get("page_type") or "").lower() == "overview"), None) or next( (page for page in pages if not str(page.get("parent_slug") or "").strip() and int(page.get("depth") or 0) == 0), None, ) or (pages[0] if pages else None) ) def _count_public_wiki_values(pages: list[dict[str, Any]], key: str, fallback: str) -> dict[str, int]: counts: dict[str, int] = {} for page in pages: value = str(page.get(key) or fallback).strip().lower() or fallback counts[value] = counts.get(value, 0) + 1 return counts def _public_wiki_time_bound(pages: list[dict[str, Any]], key: str, *, latest: bool) -> str | None: values = [str(page.get(key) or "").strip() for page in pages] values = [value for value in values if value] if not values: return None return max(values) if latest else min(values) async def _all_public_wiki_page_views( client: WeKnoraClient, remote_id: str, ) -> tuple[list[dict[str, Any]], int, bool]: pages_by_slug: dict[str, dict[str, Any]] = {} remote_total = 0 complete = False for page_number in range(1, _PUBLIC_WIKI_PAGE_LIMIT + 1): result = await client.list_wiki_pages( remote_id, page=page_number, page_size=_PUBLIC_WIKI_PAGE_SIZE, query="", ) raw_pages = result.get("pages") if isinstance(result, dict) else [] page_items = raw_pages if isinstance(raw_pages, list) else [] for page_item in page_items: if not isinstance(page_item, dict): continue page = _wiki_page_view(page_item) if page is None: continue slug = str(page.get("slug") or "").strip() pages_by_slug[slug or f"__page_{len(pages_by_slug)}"] = page reported_total = int(result.get("total") or 0) remote_total = max(remote_total, reported_total, len(pages_by_slug)) total_pages = max(1, int(result.get("total_pages") or 1)) if not page_items or page_number >= total_pages or (reported_total > 0 and len(pages_by_slug) >= remote_total): complete = True break return list(pages_by_slug.values()), max(remote_total, len(pages_by_slug)), complete def _public_wiki_summary( mapping_id: str, pages: list[dict[str, Any]], *, total: int, complete: bool, ) -> dict[str, Any]: explicit_index = next((page for page in pages if _is_public_wiki_index_page(page)), None) default_page = _public_wiki_default_page(pages) if explicit_index is not None: index_mode = "explicit" index_slug: str | None = str(explicit_index.get("slug") or "") or None elif pages: index_mode = "generated" index_slug = _GENERATED_PUBLIC_WIKI_INDEX_SLUG else: index_mode = "none" index_slug = None default_slug = index_slug or (str(default_page.get("slug") or "") if default_page else None) route = f"/embed/knowledge/{quote(mapping_id, safe='')}" if default_slug: route = f"{route}?{urlencode({'wikiSlug': default_slug})}" return { "status": "ready" if pages else "empty", "page_count": total, "loaded_page_count": len(pages), "stats_complete": complete, "page_type_counts": _count_public_wiki_values(pages, "page_type", "page"), "page_status_counts": _count_public_wiki_values(pages, "status", "unknown"), "index_mode": index_mode, "has_explicit_index": explicit_index is not None, "index_slug": index_slug, "default_slug": default_slug, "created_at": _public_wiki_time_bound(pages, "created_at", latest=False), "updated_at": _public_wiki_time_bound(pages, "updated_at", latest=True), "embed_route": route, } def _unavailable_public_wiki_summary(mapping_id: str) -> dict[str, Any]: return { "status": "unavailable", "page_count": 0, "loaded_page_count": 0, "stats_complete": False, "page_type_counts": {}, "page_status_counts": {}, "index_mode": "unknown", "has_explicit_index": False, "index_slug": None, "default_slug": None, "created_at": None, "updated_at": None, "embed_route": f"/embed/knowledge/{quote(mapping_id, safe='')}", } def _public_local_index_summary( request: Request, mapping: dict[str, Any], state: dict[str, Any] | None, ) -> dict[str, Any]: app_state = request.app.state local_config = getattr( getattr(getattr(app_state, "config", None), "llmwiki", None), "local_wiki_index", None, ) enabled = bool(getattr(local_config, "enabled", False)) current = dict(state or {}) state_name = str(current.get("state") or "not_synced") if state_name == "idle" and int(current.get("index_revision") or 0) == 0 and not current.get("last_completed_at") and not current.get("last_success_at"): state_name = "not_synced" embedding = getattr(app_state, "llmwiki_embedding", None) fingerprint = embedding.fingerprint if embedding is not None else None fully_vectorized = bool(mapping.get("wiki_index_enabled", True)) and is_fully_vectorized( current, fingerprint, ) return { "enabled": enabled, "wiki_index_enabled": bool(mapping.get("wiki_index_enabled", True)), "state": state_name, "fully_vectorized": fully_vectorized, "retrieval_mode": "vector" if fully_vectorized else "wiki_api", "remote_page_count": int(current.get("remote_page_count") or 0), "local_page_count": int(current.get("local_page_count") or 0), "ready_page_count": int(current.get("ready_page_count") or 0), "failed_page_count": int(current.get("failed_page_count") or 0), "vector_count": int(current.get("vector_count") or 0), "last_started_at": current.get("last_started_at"), "last_completed_at": current.get("last_completed_at"), "last_success_at": current.get("last_success_at"), "updated_at": current.get("updated_at"), } @router.get("/api/public/llmwiki/knowledge-bases") async def list_public_knowledge_bases(request: Request) -> dict[str, Any]: """List live, published Wiki libraries with safe metadata and index health.""" rows = [ row for row in await _store(request).list_visible( "__public_wiki__", scope="public", is_admin=False, ) if row.get("publication_status") == "published" and not _is_conversation_deposit_mapping(row) ] client = _client_or_503() try: remote_items = await client.list_knowledge_bases() except WeKnoraError as exc: _raise_upstream(exc) remote_by_id = {str(item.get("id")): item for item in remote_items if isinstance(item, dict) and item.get("id")} # A successful upstream list is authoritative: stale DeerFlow mappings are # omitted so discovery never advertises a Wiki that its detail route cannot open. rows = [row for row in rows if str(row.get("weknora_id") or "") in remote_by_id] mapping_ids = [str(row["id"]) for row in rows] index_states: dict[str, dict[str, Any]] = {} index_store = getattr(request.app.state, "llmwiki_index_store", None) if index_store is not None and mapping_ids: try: index_states = {str(state.get("knowledge_base_mapping_id") or ""): state for state in await index_store.list_index_status(mapping_ids) if isinstance(state, dict)} except Exception: logger.warning("Failed to read public Wiki local-index summaries", exc_info=True) semaphore = asyncio.Semaphore(_PUBLIC_WIKI_SUMMARY_CONCURRENCY) async def build_summary(row: dict[str, Any]) -> dict[str, Any]: mapping_id = str(row["id"]) remote_id = str(row["weknora_id"]) remote = remote_by_id[remote_id] async with semaphore: try: remote = {**remote, **await client.get_knowledge_base(remote_id)} except WeKnoraError as exc: if exc.status_code == 404: return {} logger.warning( "Failed to load public Wiki detail metadata: mapping_id=%s", mapping_id, exc_info=True, ) try: pages, page_total, complete = await _all_public_wiki_page_views(client, remote_id) wiki = _public_wiki_summary( mapping_id, pages, total=page_total, complete=complete, ) except WeKnoraError as exc: logger.warning( "Failed to load public Wiki page summary: mapping_id=%s", mapping_id, exc_info=True, ) wiki = ( _public_wiki_summary(mapping_id, [], total=0, complete=True) if exc.status_code == 404 else _unavailable_public_wiki_summary(mapping_id) ) return { **_mapping_view( row, remote, actor_user_id="__public_wiki__", is_admin=False, ), "wiki": wiki, "local_index": _public_local_index_summary( request, row, index_states.get(mapping_id), ), } knowledge_bases = [summary for summary in await asyncio.gather(*(build_summary(row) for row in rows)) if summary] return {"knowledge_bases": knowledge_bases, "total": len(knowledge_bases)} @router.get("/api/public/llmwiki/knowledge-bases/{mapping_id}") async def get_public_knowledge_base(request: Request, mapping_id: str) -> dict[str, Any]: """Return the safe public metadata used by the chrome-free Wiki reader.""" row = await _public_mapping(request, mapping_id) try: remote = await _client_or_503().get_knowledge_base(str(row["weknora_id"])) except WeKnoraError as exc: _raise_upstream(exc) return _mapping_view( row, remote, actor_user_id="__public_wiki__", is_admin=False, ) @router.get("/api/public/llmwiki/knowledge-bases/{mapping_id}/wiki/pages") async def list_public_wiki_pages( request: Request, mapping_id: str, page: int = Query(default=1, ge=1), page_size: int = Query(default=50, ge=1, le=500), query: str = Query(default="", max_length=255), ) -> dict[str, Any]: """List Wiki pages for a published knowledge base without requiring login.""" row = await _public_mapping(request, mapping_id) try: result = await _client_or_503().list_wiki_pages( str(row["weknora_id"]), page=page, page_size=page_size, query=query.strip(), ) except WeKnoraError as exc: _raise_upstream(exc) raw_pages = result.get("pages") if isinstance(result, dict) else [] pages = [view for page_item in (raw_pages if isinstance(raw_pages, list) else []) if isinstance(page_item, dict) for view in [_wiki_page_view(page_item)] if view is not None] return {**result, "pages": pages} @router.get("/api/public/llmwiki/knowledge-bases/{mapping_id}/files") async def get_public_wiki_file( request: Request, mapping_id: str, file_path: str = Query(..., min_length=1, max_length=8192), ) -> Response: """Serve one protected image belonging to a published Wiki.""" row = await _public_mapping(request, mapping_id) return await _wiki_file_response(row, file_path, public=True) @router.get("/api/public/llmwiki/files/{token}") async def get_signed_wiki_file(token: str) -> Response: """Serve a private Wiki image through a short-lived, signed HTTP URL.""" payload = _verify_signed_payload(token) if payload.get("kind") != "wiki_file": raise HTTPException(status_code=401, detail="Signed Wiki file token is invalid") mapping_id = str(payload.get("mapping_id") or "") weknora_id = str(payload.get("weknora_id") or "") file_path = str(payload.get("file_path") or "") if not mapping_id or not weknora_id or not file_path: raise HTTPException(status_code=401, detail="Signed Wiki file token is invalid") return await _wiki_file_response( {"id": mapping_id, "weknora_id": weknora_id}, file_path, public=False, ) _PUBLIC_WIKI_INDEX_TYPES = ("summary", "entity", "concept", "synthesis", "comparison") def _public_wiki_index_item(value: Any) -> dict[str, Any] | None: if not isinstance(value, dict): return None slug = str(value.get("slug") or "").strip() if not slug: return None return { "slug": slug, "title": str(value.get("title") or slug), "summary": str(value.get("summary") or ""), "category_path": ([str(item) for item in value.get("category_path", []) if str(item).strip()] if isinstance(value.get("category_path"), list) else []), "wiki_path": str(value.get("wiki_path") or ""), "depth": int(value.get("depth") or 0), } async def _all_public_wiki_index(client: WeKnoraClient, remote_id: str) -> dict[str, Any]: intro = "" version = 0 groups: list[dict[str, Any]] = [] for page_type in _PUBLIC_WIKI_INDEX_TYPES: items: list[dict[str, Any]] = [] cursor = "" seen_cursors: set[str] = set() total = 0 for _page_number in range(100): payload = await client.get_wiki_index( remote_id, types=[page_type], limit=500, cursor=cursor, ) intro = intro or str(payload.get("intro") or "") version = max(version, int(payload.get("version") or 0)) raw_groups = payload.get("groups") if isinstance(payload.get("groups"), list) else [] group = next( (value for value in raw_groups if isinstance(value, dict) and value.get("type") == page_type), {}, ) total = max(total, int(group.get("total") or 0)) raw_items = group.get("items") if isinstance(group.get("items"), list) else [] items.extend(item for value in raw_items if (item := _public_wiki_index_item(value)) is not None) next_cursor = str(group.get("next_cursor") or "") if not next_cursor or next_cursor in seen_cursors: break seen_cursors.add(next_cursor) cursor = next_cursor else: logger.warning( "Stopped public Wiki index pagination at safety limit: mapping_remote_id=%s page_type=%s", remote_id, page_type, ) groups.append({"type": page_type, "total": total or len(items), "items": items}) return {"intro": intro, "version": version, "groups": groups} @router.get("/api/public/llmwiki/knowledge-bases/{mapping_id}/wiki/index") async def get_public_wiki_index(request: Request, mapping_id: str) -> dict[str, Any]: """Return WeKnora's ordered dynamic Wiki index for the public reader.""" row = await _public_mapping(request, mapping_id) try: return await _all_public_wiki_index(_client_or_503(), str(row["weknora_id"])) except WeKnoraError as exc: _raise_upstream(exc) @router.get("/api/public/llmwiki/knowledge-bases/{mapping_id}/wiki/pages/{slug:path}") async def get_public_wiki_page( request: Request, mapping_id: str, slug: str, ) -> dict[str, Any]: """Read one Wiki article from a published knowledge base.""" row = await _public_mapping(request, mapping_id) try: page = await _client_or_503().get_wiki_page(str(row["weknora_id"]), slug) except WeKnoraError as exc: _raise_upstream(exc) return _wiki_page_view(page) or page @router.get("/api/llmwiki/knowledge-bases") async def list_knowledge_bases( request: Request, scope: Literal["personal", "public", "all"] = Query(default="all"), ) -> dict[str, Any]: client = _client_or_503() user_id, is_admin = await _actor(request) try: remote_items = await client.list_knowledge_bases() remote = {str(item.get("id")): item for item in remote_items if item.get("id")} except WeKnoraError as exc: _raise_upstream(exc) deposit_remote = next( (item for item in remote_items if str(item.get("name") or "").strip() == CONVERSATION_DEPOSIT_KB_NAME), None, ) if deposit_remote is not None: previous_mapping = await _store(request).get_by_weknora_id(str(deposit_remote.get("id") or "")) deposit_mapping = await ensure_conversation_deposit_mapping( _store(request), client, remote_knowledge_bases=remote_items, ) description_was_synchronized = previous_mapping is None or not _is_conversation_deposit_mapping(previous_mapping) or str(previous_mapping.get("description") or "") != CONVERSATION_DEPOSIT_KB_DESCRIPTION if deposit_mapping is not None and description_was_synchronized: remote_id = str(deposit_mapping["weknora_id"]) remote[remote_id] = { **remote.get(remote_id, deposit_remote), "description": CONVERSATION_DEPOSIT_KB_DESCRIPTION, } rows = await _store(request).list_visible(user_id, scope=scope, is_admin=False) knowledge_bases: list[dict[str, Any]] = [] for row in rows: remote_item = remote.get(str(row["weknora_id"])) # A successful remote list is authoritative for every scope. Keep the # mapping for existing references, but never display a deleted base. if remote_item is None: continue row, remote_item = await _sync_conversation_deposit_view(_store(request), client, row, remote_item) knowledge_bases.append(_mapping_view(row, remote_item, actor_user_id=user_id, is_admin=is_admin)) return {"knowledge_bases": knowledge_bases} @router.get("/api/llmwiki/agents/{agent_id}/knowledge-bases") async def list_agent_knowledge_bases(request: Request, agent_id: str) -> dict[str, Any]: """Return the agent-bound bases still visible to the current DeerFlow user.""" user_id, is_admin = await _actor(request) agent_store = get_agent_store(request) agent = await (agent_store.get_any(agent_id) if is_admin else agent_store.get_visible(agent_id, user_id)) if agent is None: raise HTTPException(status_code=404, detail="Agent not found") try: config = load_agent_config(agent_id) except (FileNotFoundError, ValueError): config = None bound_ids = list(config.llmwiki_knowledge_base_ids or []) if config else [] rows: list[dict[str, Any]] = [] for mapping_id in bound_ids: row = await _store(request).get_authorized( mapping_id, user_id, write=False, is_admin=is_admin, ) if row is not None: rows.append(row) try: remote = {str(item.get("id")): item for item in await _client_or_503().list_knowledge_bases() if item.get("id")} except WeKnoraError as exc: _raise_upstream(exc) return {"knowledge_bases": [_mapping_view(row, remote_item, actor_user_id=user_id, is_admin=is_admin) for row in rows if (remote_item := remote.get(str(row["weknora_id"]))) is not None]} @router.post("/api/llmwiki/knowledge-bases", status_code=201) async def create_knowledge_base(request: Request, body: KnowledgeBaseCreate) -> dict[str, Any]: client = _client_or_503() user_id, is_admin = await _actor(request) if body.name.strip() == CONVERSATION_DEPOSIT_KB_NAME: raise HTTPException(status_code=409, detail="对话沉淀是 cmzs 系统保留知识库") try: remote = await client.create_knowledge_base( name=body.name.strip(), description=body.description, kb_type=body.type, wiki_enabled=body.wiki_enabled, ) except WeKnoraError as exc: _raise_upstream(exc) try: row = await _store(request).create_mapping( weknora_id=str(remote["id"]), owner_user_id=user_id, name=str(remote.get("name") or body.name), description=str(remote.get("description") or body.description), kb_type=str(remote.get("type") or body.type), ) except Exception: logger.exception("Failed to persist LLMWiki mapping; compensating remote create") try: await client.delete_knowledge_base(str(remote["id"])) except Exception: logger.exception("Failed to compensate orphan WeKnora knowledge base") raise HTTPException(status_code=500, detail="Failed to register the knowledge base") from None _schedule_local_wiki_sync(request, row) return _mapping_view(row, remote, actor_user_id=user_id, is_admin=is_admin) @router.get("/api/llmwiki/knowledge-bases/{mapping_id}") async def get_knowledge_base(request: Request, mapping_id: str) -> dict[str, Any]: row, user_id, is_admin = await _authorized_mapping(request, mapping_id, write=False) client = _client_or_503() try: remote = await client.get_knowledge_base(str(row["weknora_id"])) except WeKnoraError as exc: _raise_upstream(exc) row, remote = await _sync_conversation_deposit_view(_store(request), client, row, remote) return _mapping_view(row, remote, actor_user_id=user_id, is_admin=is_admin) @router.patch("/api/llmwiki/knowledge-bases/{mapping_id}") async def update_knowledge_base(request: Request, mapping_id: str, body: KnowledgeBaseUpdate) -> dict[str, Any]: row, user_id, is_admin = await _authorized_mapping(request, mapping_id, write=True) changes = body.model_dump(exclude_none=True) if "name" in changes: changes["name"] = changes["name"].strip() if not changes["name"]: raise HTTPException(status_code=422, detail="知识库名称不能为空") if changes["name"] == CONVERSATION_DEPOSIT_KB_NAME: raise HTTPException(status_code=409, detail="对话沉淀是 cmzs 系统保留知识库") if not changes: return _mapping_view(row, None, actor_user_id=user_id, is_admin=is_admin) try: remote = await _client_or_503().update_knowledge_base(str(row["weknora_id"]), changes) except WeKnoraError as exc: _raise_upstream(exc) updated = await _store(request).update_mapping(mapping_id, name=changes.get("name"), description=changes.get("description")) return _mapping_view(updated or row, remote, actor_user_id=user_id, is_admin=is_admin) @router.delete("/api/llmwiki/knowledge-bases/{mapping_id}", status_code=204) async def delete_knowledge_base(request: Request, mapping_id: str) -> None: row, _, _ = await _authorized_mapping(request, mapping_id, write=True) try: await _client_or_503().delete_knowledge_base(str(row["weknora_id"])) except WeKnoraError as exc: # A 404 means the remote data is already gone. The owner/admin must still # be able to remove DeerFlow's stale mapping; other upstream failures stay # visible and never cause local-only deletion. if exc.status_code != 404: _raise_upstream(exc) await _store(request).delete_mapping(mapping_id) @router.get("/api/llmwiki/knowledge-bases/{mapping_id}/documents") async def list_documents( request: Request, mapping_id: str, page: int = Query(default=1, ge=1), page_size: int = Query(default=100, ge=1, le=200), ) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, mapping_id, write=False) try: result = await _client_or_503().list_documents(str(row["weknora_id"]), page=page, page_size=page_size) except WeKnoraError as exc: _raise_upstream(exc) return {**result, "items": [_document_view(item) for item in result["items"]]} @router.get("/api/llmwiki/knowledge-bases/{mapping_id}/documents/{document_id}") async def get_document(request: Request, mapping_id: str, document_id: str) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, mapping_id, write=False) try: document = await _verify_document(_client_or_503(), row, document_id) except WeKnoraError as exc: _raise_upstream(exc) return _document_view(document) @router.get("/api/llmwiki/knowledge-bases/{mapping_id}/documents/{document_id}/chunks") async def list_document_chunks( request: Request, mapping_id: str, document_id: str, page: int = Query(default=1, ge=1), page_size: int = Query(default=20, ge=1, le=100), ) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, mapping_id, write=False) client = _client_or_503() try: await _verify_document(client, row, document_id) result = await client.list_chunks(document_id, page=page, page_size=page_size) except WeKnoraError as exc: _raise_upstream(exc) return {**result, "items": [_chunk_view(item) for item in result["items"]]} @router.get("/api/llmwiki/knowledge-bases/{mapping_id}/documents/{document_id}/preview") async def preview_document(request: Request, mapping_id: str, document_id: str) -> Response: row, _, _ = await _authorized_mapping(request, mapping_id, write=False) client = _client_or_503() try: document = await _verify_document(client, row, document_id) content, media_type = await client.get_document_preview(document_id) except WeKnoraError as exc: _raise_upstream(exc) filename = str(document.get("file_name") or document.get("title") or "document") return Response( content=content, media_type=media_type, headers={ "Content-Disposition": f"inline; filename*=UTF-8''{quote(filename)}", "Cache-Control": "private, max-age=60", }, ) @router.get("/api/llmwiki/knowledge-bases/{mapping_id}/files") async def get_wiki_file( request: Request, mapping_id: str, file_path: str = Query(..., min_length=1, max_length=8192), ) -> Response: """Serve one protected Wiki image after normal knowledge-base authorization.""" row, _, _ = await _authorized_mapping(request, mapping_id, write=False) return await _wiki_file_response(row, file_path, public=False) @router.get("/api/llmwiki/knowledge-bases/{mapping_id}/files/access-url") async def create_wiki_file_access( request: Request, mapping_id: str, file_path: str = Query(..., min_length=1, max_length=8192), ) -> dict[str, Any]: """Issue a browser-usable HTTP capability after normal KB authorization.""" row, _, _ = await _authorized_mapping(request, mapping_id, write=False) safe_file_path = _validate_protected_wiki_file_path(file_path) token = _sign_embed_payload( { "kind": "wiki_file", "mapping_id": mapping_id, "weknora_id": str(row["weknora_id"]), "file_path": safe_file_path, "exp": int(time.time()) + _WIKI_FILE_ACCESS_TTL_SECONDS, } ) return {"token": token, "expires_in": _WIKI_FILE_ACCESS_TTL_SECONDS} @router.get("/api/llmwiki/knowledge-bases/{mapping_id}/graph") async def get_knowledge_base_graph( request: Request, mapping_id: str, limit: int = Query(default=200, ge=1, le=500), ) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, mapping_id, write=False) client = _client_or_503() try: remote = await client.get_knowledge_base(str(row["weknora_id"])) strategy = remote.get("indexing_strategy") if isinstance(remote.get("indexing_strategy"), dict) else {} if not strategy.get("wiki_enabled"): return { "nodes": [], "edges": [], "meta": {"mode": "disabled", "total": 0, "returned": 0, "truncated": False}, } graph = await client.get_wiki_graph(str(row["weknora_id"]), limit=limit) except WeKnoraError as exc: _raise_upstream(exc) nodes = [ { "id": str(node.get("slug") or ""), "title": str(node.get("title") or node.get("slug") or ""), "type": str(node.get("page_type") or "page"), "link_count": int(node.get("link_count") or 0), } for node in graph.get("nodes", []) if isinstance(node, dict) and node.get("slug") ] edges = [{"source": str(edge.get("source") or ""), "target": str(edge.get("target") or "")} for edge in graph.get("edges", []) if isinstance(edge, dict) and edge.get("source") and edge.get("target")] meta = graph.get("meta") if isinstance(graph.get("meta"), dict) else {} return {"nodes": nodes, "edges": edges, "meta": meta} @router.get("/api/llmwiki/knowledge-bases/{mapping_id}/wiki/pages") async def list_wiki_pages( request: Request, mapping_id: str, page: int = Query(default=1, ge=1), page_size: int = Query(default=50, ge=1, le=500), query: str = Query(default="", max_length=255), ) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, mapping_id, write=False) try: result = await _client_or_503().list_wiki_pages( str(row["weknora_id"]), page=page, page_size=page_size, query=query.strip(), ) raw_pages = result.get("pages") if isinstance(result, dict) else [] pages = [view for page_item in (raw_pages if isinstance(raw_pages, list) else []) if isinstance(page_item, dict) for view in [_wiki_page_view(page_item)] if view is not None] return {**result, "pages": pages} except WeKnoraError as exc: _raise_upstream(exc) @router.get("/api/llmwiki/knowledge-bases/{mapping_id}/wiki/export.xlsx") async def export_wiki_pages_excel(request: Request, mapping_id: str) -> StreamingResponse: """Export every processed Wiki page with its current and parent directories.""" row, _, _ = await _authorized_mapping(request, mapping_id, write=False) client = _client_or_503() try: pages = await _all_wiki_pages(client, str(row["weknora_id"])) except WeKnoraError as exc: _raise_upstream(exc) try: output = await asyncio.to_thread(_build_wiki_excel, str(row.get("name") or "LLMWiki"), pages) except ImportError as exc: raise HTTPException(status_code=500, detail=f"openpyxl is not installed: {exc}") from exc timestamp = datetime.now(UTC).strftime("%Y%m%d-%H%M%S") safe_name = re.sub(r"[\\/:*?\"<>|\r\n]+", "_", str(row.get("name") or "LLMWiki")).strip(" ._") or "LLMWiki" filename = f"{safe_name}-Wiki-{timestamp}.xlsx" return StreamingResponse( output, media_type="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet", headers={"Content-Disposition": f"attachment; filename=\"wiki-export.xlsx\"; filename*=UTF-8''{quote(filename)}"}, ) @router.post("/api/llmwiki/knowledge-bases/{mapping_id}/wiki/pages") async def create_wiki_page(request: Request, mapping_id: str, body: WikiPageCreate) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, mapping_id, write=True) payload = body.model_dump(exclude_none=True) try: page = await _client_or_503().create_wiki_page(str(row["weknora_id"]), payload) _schedule_local_wiki_sync(request, row) return _wiki_page_view(page) or page except WeKnoraError as exc: _raise_upstream(exc) @router.get("/api/llmwiki/knowledge-bases/{mapping_id}/wiki/pages/{slug:path}") async def get_wiki_page(request: Request, mapping_id: str, slug: str) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, mapping_id, write=False) try: page = await _client_or_503().get_wiki_page(str(row["weknora_id"]), slug) return _wiki_page_view(page) or page except WeKnoraError as exc: _raise_upstream(exc) @router.patch("/api/llmwiki/knowledge-bases/{mapping_id}/wiki/pages/{slug:path}") async def update_wiki_page(request: Request, mapping_id: str, slug: str, body: WikiPageUpdate) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, mapping_id, write=True) changes = body.model_dump(exclude_none=True) if not changes: try: page = await _client_or_503().get_wiki_page(str(row["weknora_id"]), slug) return _wiki_page_view(page) or page except WeKnoraError as exc: _raise_upstream(exc) try: page = await _client_or_503().update_wiki_page(str(row["weknora_id"]), slug, changes) _schedule_local_wiki_sync(request, row) return _wiki_page_view(page) or page except WeKnoraError as exc: _raise_upstream(exc) @router.delete("/api/llmwiki/knowledge-bases/{mapping_id}/wiki/pages/{slug:path}", status_code=204) async def delete_wiki_page(request: Request, mapping_id: str, slug: str) -> None: row, _, _ = await _authorized_mapping(request, mapping_id, write=True) try: await _client_or_503().delete_wiki_page(str(row["weknora_id"]), slug) _schedule_local_wiki_sync(request, row) except WeKnoraError as exc: _raise_upstream(exc) @router.post("/api/llmwiki/knowledge-bases/{mapping_id}/documents", status_code=201) async def upload_document(request: Request, mapping_id: str, file: UploadFile = File(...)) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, mapping_id, write=True) content = await file.read(_MAX_UPLOAD_BYTES + 1) if len(content) > _MAX_UPLOAD_BYTES: raise HTTPException(status_code=413, detail="File exceeds the 100 MB cmzs upload limit") if not file.filename: raise HTTPException(status_code=422, detail="File name is required") try: result = await _client_or_503().upload_document( str(row["weknora_id"]), filename=file.filename, content=content, content_type=file.content_type or "application/octet-stream", ) _schedule_local_wiki_sync(request, row) return result except WeKnoraError as exc: _raise_upstream(exc) @router.post("/api/llmwiki/knowledge-bases/{mapping_id}/documents/url", status_code=201) async def import_document_url(request: Request, mapping_id: str, body: KnowledgeUrlImport) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, mapping_id, write=True) url = body.url.strip() if not url.lower().startswith(("http://", "https://")): raise HTTPException(status_code=422, detail="URL must start with http:// or https://") try: result = await _client_or_503().import_document_url(str(row["weknora_id"]), url=url) _schedule_local_wiki_sync(request, row) return result except WeKnoraError as exc: _raise_upstream(exc) @router.post("/api/llmwiki/knowledge-bases/{mapping_id}/documents/manual", status_code=201) async def create_manual_document(request: Request, mapping_id: str, body: ManualKnowledgeCreate) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, mapping_id, write=True) try: result = await _client_or_503().create_manual_document( str(row["weknora_id"]), title=body.title.strip(), content=body.content, ) _schedule_local_wiki_sync(request, row) return result except WeKnoraError as exc: _raise_upstream(exc) @router.delete("/api/llmwiki/knowledge-bases/{mapping_id}/documents/{document_id}", status_code=204) async def delete_document(request: Request, mapping_id: str, document_id: str) -> None: row, _, _ = await _authorized_mapping(request, mapping_id, write=True) client = _client_or_503() try: await _verify_document(client, row, document_id) await client.delete_document(document_id) _schedule_local_wiki_sync(request, row) except WeKnoraError as exc: _raise_upstream(exc) @router.post("/api/llmwiki/knowledge-bases/{mapping_id}/documents/{document_id}/reprocess") async def reprocess_document(request: Request, mapping_id: str, document_id: str) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, mapping_id, write=True) client = _client_or_503() try: await _verify_document(client, row, document_id) result = await client.reprocess_document(document_id) _schedule_local_wiki_sync(request, row) return result except WeKnoraError as exc: _raise_upstream(exc) @router.post("/api/llmwiki/search") async def search_llmwiki(request: Request, body: SearchRequest) -> dict[str, Any]: user_id, is_admin = await _actor(request) store = _store(request) if body.knowledge_base_ids: rows: list[dict[str, Any]] = [] for mapping_id in dict.fromkeys(body.knowledge_base_ids): row = await store.get_authorized(mapping_id, user_id, write=False, is_admin=is_admin) if row is None: raise HTTPException(status_code=404, detail="Knowledge base not found") rows.append(row) else: rows = await store.list_visible(user_id, scope=body.scope, is_admin=False) service = getattr(request.app.state, "llmwiki_retrieval_service", None) if service is None: raise HTTPException(status_code=503, detail={"code": "WIKI_RETRIEVAL_UNAVAILABLE"}) include_drafts = body.include_drafts and (is_admin or all(str(row.get("owner_user_id")) == user_id for row in rows)) try: result = await service.search(body.query, rows, top_k_pages=body.top_k, include_drafts=include_drafts) except WikiIndexError as exc: raise HTTPException(status_code=exc.status_code, detail={"code": exc.code, "message": str(exc)}) from None mapping_by_id = {str(row["id"]): row for row in rows} result["results"] = [_sanitize_conversation_deposit_value(item) if _is_conversation_deposit_mapping(mapping_by_id.get(str(item.get("knowledge_base_id") or ""), {})) else item for item in result.get("results") or []] return result @router.post("/api/llmwiki/conversation-deposits/turn", status_code=202) async def create_conversation_deposit( request: Request, background_tasks: BackgroundTasks, body: ConversationDepositCreate, ) -> dict[str, Any]: runtime = get_resolved_llmwiki_runtime(getattr(request.app.state, "config", None)) if not runtime.weknora_enabled: return {"queued": False, "provider": runtime.provider} user_id, _ = await _actor(request) if not await get_thread_store(request).check_access(body.thread_id, user_id, require_existing=True): raise HTTPException(status_code=404, detail="Thread not found") app = request.app turn = ConversationTurnDeposit( thread_id=body.thread_id, question=body.question, answer=body.answer, human_message_id=body.human_message_id, assistant_message_id=body.assistant_message_id, assistant_id=body.assistant_id, knowledge_base_ids=list(dict.fromkeys(item.strip() for item in body.knowledge_base_ids if item.strip())), created_by_user_id=user_id, ) background_tasks.add_task(deposit_conversation_turn, app, turn) return {"queued": True, "provider": "weknora"} @router.post("/api/llmwiki/knowledge-bases/{mapping_id}/publish") async def request_publish(request: Request, mapping_id: str) -> dict[str, Any]: user_id, is_admin = await _actor(request) existing = await _store(request).get_authorized( mapping_id, user_id, write=True, is_admin=is_admin, ) if existing is None: raise HTTPException(status_code=404, detail="Knowledge base not found") if _is_conversation_deposit_mapping(existing) and existing.get("publication_status") == "published": return {"id": existing["id"], "publication_status": "published"} try: await _client_or_503().get_knowledge_base(str(existing["weknora_id"])) except WeKnoraError as exc: _raise_upstream(exc) row = await _store(request).request_publish(mapping_id, user_id, is_admin=is_admin) if row is None: raise HTTPException(status_code=404, detail="Knowledge base not found") return {"id": row["id"], "publication_status": row["publication_status"]} @router.post("/api/llmwiki/knowledge-bases/{mapping_id}/unpublish") async def unpublish(request: Request, mapping_id: str) -> dict[str, Any]: user_id, is_admin = await _actor(request) existing = await _store(request).get_authorized(mapping_id, user_id, write=True, is_admin=is_admin) if existing is None: raise HTTPException(status_code=404, detail="Knowledge base not found") if _is_conversation_deposit_mapping(existing): raise HTTPException(status_code=403, detail="Conversation deposit is always public") row = await _store(request).unpublish(mapping_id, user_id, is_admin=is_admin) if row is None: raise HTTPException(status_code=404, detail="Knowledge base not found") return {"id": row["id"], "publication_status": row["publication_status"]} @router.get("/api/llmwiki/sources/context") async def get_source_context( request: Request, knowledge_base_id: str = Query(..., min_length=1, max_length=128), chunk_id: str = Query(..., min_length=1, max_length=128), ) -> dict[str, Any]: row, _, _ = await _authorized_mapping(request, knowledge_base_id, write=False) client = _client_or_503() try: chunks = await client.get_chunk_context( chunk_id, expected_knowledge_base_id=str(row["weknora_id"]), radius=2, ) except WeKnoraError as exc: _raise_upstream(exc) wiki_page = await _find_wiki_page_for_chunk( client, knowledge_base_id=str(row["weknora_id"]), chunk_id=chunk_id, chunks=chunks, ) sanitize_deposit = _is_conversation_deposit_mapping(row) if sanitize_deposit and wiki_page is not None: wiki_page = _sanitize_conversation_deposit_value(wiki_page) focus_chunk = next((chunk for chunk in chunks if str(chunk.get("id") or "") == chunk_id), {}) document_id = str(focus_chunk.get("knowledge_id") or focus_chunk.get("document_id") or focus_chunk.get("knowledgeId") or "") document_title = str(focus_chunk.get("knowledge_title") or focus_chunk.get("title") or focus_chunk.get("file_name") or focus_chunk.get("filename") or "") normalized_chunks = [ { "id": str(chunk.get("id") or ""), "chunk_index": int(chunk.get("chunk_index") or 0), "content": str(_sanitize_conversation_deposit_value(chunk.get("content") or "") if sanitize_deposit else chunk.get("content") or ""), "is_focus": str(chunk.get("id") or "") == chunk_id, } for chunk in chunks if chunk.get("id") ] return { "success": True, "data": { "knowledge_base_id": row["id"], "knowledge_base_name": row.get("name") or "", "document_id": _sanitize_conversation_deposit_value(document_id) if sanitize_deposit else document_id, "document_title": _sanitize_conversation_deposit_value(document_title) if sanitize_deposit else document_title, "focus_chunk_id": chunk_id, "display_mode": "wiki" if wiki_page is not None else "chunks", "wiki_page": wiki_page, "chunks": normalized_chunks, }, }