"""File-backed assistant knowledge import/export packages. The assistant knowledge base is fed by WeKnora-exported data. Keep that boundary explicit: a WeKnora export is first materialized as a durable JSONL package, then imported asynchronously from that package. JSONL keeps large exports append-friendly and avoids treating the HTTP request body as the job state. """ from __future__ import annotations import json from collections.abc import Iterable, Mapping from pathlib import Path from typing import Any PACKAGE_VERSION = 1 PACKAGE_RECORD_TYPES = { "source", "wiki_page", "document", "chunk", "vector", "entity", "relation", "graph", } def package_path(root: Path, job_id: str) -> Path: safe = "".join(ch for ch in str(job_id) if ch.isalnum() or ch in {"-", "_"}) return root / f"{safe}.jsonl" def write_jsonl_record(handle: Any, record_type: str, payload: Mapping[str, Any]) -> None: handle.write(json.dumps({"record_type": record_type, "payload": dict(payload)}, ensure_ascii=False)) handle.write("\n") def write_package(path: Path, records: Iterable[tuple[str, Mapping[str, Any]]]) -> dict[str, int]: path.parent.mkdir(parents=True, exist_ok=True) counts: dict[str, int] = {key: 0 for key in PACKAGE_RECORD_TYPES} tmp_path = path.with_suffix(path.suffix + ".tmp") with tmp_path.open("w", encoding="utf-8", newline="\n") as handle: write_jsonl_record( handle, "source", { "package_version": PACKAGE_VERSION, }, ) counts["source"] += 1 for record_type, payload in records: if record_type not in PACKAGE_RECORD_TYPES: continue write_jsonl_record(handle, record_type, payload) counts[record_type] = counts.get(record_type, 0) + 1 tmp_path.replace(path) return counts def read_package(path: Path) -> dict[str, Any]: package: dict[str, Any] = { "source": {}, "wiki_pages": [], "documents": [], "chunks": [], "vectors": [], "entities": [], "relations": [], "graph": {"nodes": [], "edges": []}, } with path.open("r", encoding="utf-8") as handle: for line in handle: line = line.strip() if not line: continue record = json.loads(line) if not isinstance(record, dict): continue record_type = str(record.get("record_type") or "") payload = record.get("payload") if not isinstance(payload, dict): continue if record_type == "source": package["source"] = {**package.get("source", {}), **payload} elif record_type == "wiki_page": package["wiki_pages"].append(payload) elif record_type == "document": package["documents"].append(payload) elif record_type == "chunk": package["chunks"].append(payload) elif record_type == "vector": package["vectors"].append(payload) elif record_type == "entity": package["entities"].append(payload) elif record_type == "relation": package["relations"].append(payload) elif record_type == "graph": nodes = payload.get("nodes") edges = payload.get("edges") if isinstance(nodes, list): package["graph"]["nodes"].extend(item for item in nodes if isinstance(item, dict)) if isinstance(edges, list): package["graph"]["edges"].extend(item for item in edges if isinstance(item, dict)) for key, value in payload.items(): if key not in {"nodes", "edges"}: package["graph"][key] = value return package