themis / phase1 /scripts /knowledge_service.py
vg15o2's picture
Moonley backend (HF Space build)
1d9bd9b
Raw
History Blame Contribute Delete
11.4 kB
"""Private project-document extraction and optional Qdrant vector publishing."""
from __future__ import annotations
import hashlib
import json
import os
import re
import uuid
from pathlib import Path
from urllib.parse import quote
import numpy as np
import requests
from document_text import extract_document
from project_store import ProjectStore
class KnowledgeServiceError(Exception):
pass
def _chunks(text: str, *, size: int = 1_400, overlap: int = 180, limit: int = 64) -> list[str]:
paragraphs = [re.sub(r"\s+", " ", value).strip() for value in re.split(r"\n\s*\n", text) if value.strip()]
chunks: list[str] = []
current = ""
for paragraph in paragraphs:
pending = paragraph
while pending:
room = size - len(current)
if room <= 80:
chunks.append(current.strip())
current = current[-overlap:].lstrip()
room = size - len(current)
take = pending[:room]
split = take.rfind(" ") if len(pending) > room else len(take)
if split < max(80, room // 2):
split = len(take)
current = (current + " " + pending[:split]).strip()
pending = pending[split:].lstrip()
if len(chunks) >= limit:
return chunks[:limit]
if current and len(chunks) < limit:
chunks.append(current.strip())
return chunks[:limit]
class KnowledgeService:
def __init__(self, projects: ProjectStore, corpus):
self.projects = projects
self.corpus = corpus
self.qdrant_url = os.environ.get("QDRANT_URL", "").rstrip("/")
self.qdrant_key = os.environ.get("QDRANT_API_KEY", "")
self.collection = os.environ.get("QDRANT_KNOWLEDGE_COLLECTION", "moonley_tenant_knowledge")
self.supabase_url = os.environ.get("SUPABASE_URL", "").rstrip("/")
self.supabase_key = os.environ.get("SUPABASE_SERVICE_ROLE_KEY", "")
self.supabase_bucket = os.environ.get("SUPABASE_STORAGE_BUCKET", "")
@property
def qdrant_configured(self) -> bool:
return bool(self.qdrant_url and self.qdrant_key)
@property
def supabase_configured(self) -> bool:
return bool(self.supabase_url and self.supabase_key and self.supabase_bucket)
def status(self) -> dict:
return {
"source_provider": "supabase_private" if self.supabase_configured else ("mounted_volume" if self.projects.configured else "unconfigured"),
"vector_provider": "qdrant" if self.qdrant_configured else "private_mounted_volume",
"qdrant_configured": self.qdrant_configured,
"supabase_configured": self.supabase_configured,
"ocr": "tesseract",
"chat_embedding_default": False,
}
@staticmethod
def _owner_key(owner_id: str) -> str:
return hashlib.sha256(owner_id.encode("utf-8")).hexdigest()
def source_text(self, owner_id: str, project_id: str, document_id: str) -> tuple[str, dict]:
knowledge_dir = self.projects.knowledge_dir(owner_id, project_id, document_id)
cache_path = knowledge_dir / "content.json"
if cache_path.exists():
payload = json.loads(cache_path.read_text(encoding="utf-8"))
return str(payload.get("text") or ""), dict(payload.get("extraction") or {})
document = self.projects.document_record(owner_id, project_id, document_id)
result = extract_document(
self.projects.document_path(owner_id, project_id, document_id),
str(document.get("media_type") or ""),
)
extraction = result.public_dict()
temp = knowledge_dir / f".content-{uuid.uuid4().hex}.tmp"
temp.write_text(
json.dumps({"version": 1, "text": result.text, "extraction": extraction}, ensure_ascii=False),
encoding="utf-8",
)
os.replace(temp, cache_path)
self.projects.record_extraction(owner_id, project_id, document_id, extraction, status="extracted")
return result.text, extraction
def ingest(self, owner_id: str, project_id: str, document_id: str) -> None:
try:
if self.supabase_configured:
self._publish_supabase(owner_id, project_id, document_id)
text, extraction = self.source_text(owner_id, project_id, document_id)
chunks = _chunks(text)
if not chunks:
raise ValueError("No readable text was found in the document.")
vectors = self.corpus.encode_documents(chunks)
knowledge_dir = self.projects.knowledge_dir(owner_id, project_id, document_id)
chunk_payload = {
"version": 1,
"chunks": [{"id": index, "text": value} for index, value in enumerate(chunks)],
}
chunks_temp = knowledge_dir / f".chunks-{uuid.uuid4().hex}.tmp"
chunks_temp.write_text(json.dumps(chunk_payload, ensure_ascii=False), encoding="utf-8")
os.replace(chunks_temp, knowledge_dir / "chunks.json")
provider = "qdrant" if self.qdrant_configured else "private_mounted_volume"
if self.qdrant_configured:
self._publish_qdrant(owner_id, project_id, document_id, vectors)
else:
vector_temp = knowledge_dir / f".vectors-{uuid.uuid4().hex}.npy"
np.save(vector_temp, vectors)
os.replace(vector_temp, knowledge_dir / "vectors.npy")
self.projects.record_extraction(
owner_id,
project_id,
document_id,
{
**extraction,
"chunk_count": len(chunks),
"vector_provider": provider,
"source_provider": "supabase_private" if self.supabase_configured else "mounted_volume",
},
status="ready",
)
except Exception as exc:
try:
self.projects.record_extraction(
owner_id,
project_id,
document_id,
{"method": "failed", "text_chars": 0},
status="failed",
)
except Exception:
pass
print(f"[knowledge] ingestion failed document={document_id}: {type(exc).__name__}", flush=True)
def _supabase_headers(self, media_type: str | None = None) -> dict[str, str]:
headers = {
"apikey": self.supabase_key,
"Authorization": f"Bearer {self.supabase_key}",
"x-upsert": "true",
}
if media_type:
headers["Content-Type"] = media_type
return headers
def _source_object_path(self, owner_id: str, project_id: str, document: dict) -> str:
suffix = Path(str(document.get("stored_name") or "")).suffix.lower()
return f"users/{self._owner_key(owner_id)}/projects/{project_id}/documents/{document['id']}{suffix}"
def _publish_supabase(self, owner_id: str, project_id: str, document_id: str) -> None:
document = self.projects.document_record(owner_id, project_id, document_id)
object_path = self._source_object_path(owner_id, project_id, document)
endpoint = (
f"{self.supabase_url}/storage/v1/object/{quote(self.supabase_bucket, safe='')}/"
f"{quote(object_path, safe='/')}"
)
response = requests.post(
endpoint,
headers=self._supabase_headers(str(document.get("media_type") or "application/octet-stream")),
data=self.projects.document_path(owner_id, project_id, document_id).read_bytes(),
timeout=90,
)
if response.status_code not in {200, 201}:
raise KnowledgeServiceError(f"Private Supabase upload failed ({response.status_code}).")
def delete(self, owner_id: str, project_id: str, document_id: str) -> None:
document = self.projects.document_record(owner_id, project_id, document_id)
if self.qdrant_configured:
tenant_id = self._owner_key(owner_id)
response = requests.post(
f"{self.qdrant_url}/collections/{self.collection}/points/delete?wait=true",
headers=self._headers(),
timeout=45,
json={
"filter": {
"must": [
{"key": "tenant_id", "match": {"value": tenant_id}},
{"key": "project_id", "match": {"value": project_id}},
{"key": "document_id", "match": {"value": document_id}},
]
}
},
)
if response.status_code not in {200, 404}:
raise KnowledgeServiceError(f"Qdrant deletion failed ({response.status_code}).")
if self.supabase_configured:
object_path = self._source_object_path(owner_id, project_id, document)
endpoint = (
f"{self.supabase_url}/storage/v1/object/{quote(self.supabase_bucket, safe='')}/"
f"{quote(object_path, safe='/')}"
)
response = requests.delete(endpoint, headers=self._supabase_headers(), timeout=45)
if response.status_code not in {200, 404}:
raise KnowledgeServiceError(f"Private Supabase deletion failed ({response.status_code}).")
def delete_project(self, owner_id: str, project_id: str) -> None:
project = self.projects.get_project(owner_id, project_id)
for document in project.get("documents") or []:
self.delete(owner_id, project_id, str(document.get("id") or ""))
def _headers(self) -> dict[str, str]:
return {"api-key": self.qdrant_key, "Content-Type": "application/json"}
def _publish_qdrant(self, owner_id: str, project_id: str, document_id: str, vectors: np.ndarray) -> None:
collection_url = f"{self.qdrant_url}/collections/{self.collection}"
response = requests.get(collection_url, headers=self._headers(), timeout=15)
if response.status_code == 404:
response = requests.put(
collection_url,
headers=self._headers(),
timeout=30,
json={"vectors": {"size": int(vectors.shape[1]), "distance": "Cosine"}, "on_disk_payload": True},
)
response.raise_for_status()
tenant_id = self._owner_key(owner_id)
namespace = uuid.UUID("db853e8b-aeb1-47c8-a7fc-680c662ba8ee")
points = [
{
"id": str(uuid.uuid5(namespace, f"{tenant_id}:{project_id}:{document_id}:{index}")),
"vector": vector.tolist(),
"payload": {
"tenant_id": tenant_id,
"project_id": project_id,
"document_id": document_id,
"chunk_id": index,
},
}
for index, vector in enumerate(vectors)
]
for start in range(0, len(points), 32):
result = requests.put(
f"{collection_url}/points?wait=true",
headers=self._headers(),
timeout=90,
json={"points": points[start:start + 32]},
)
result.raise_for_status()