Files
microfish/backend/app/api/graph.py
Kunthawat Greethong 8b84378fe1 feat: SaaS foundation for CrowdSight
Elevate MiroFish/CrowdSight from single-container dev to a SaaS foundation:

- Local memory backend (Zep-compatible): memory services/models, local graph
  builder + updater, AgentActivity seam, import-boundary isolation; Zep stays
  default, local is opt-in behind MEMORY_BACKEND. Semantic parity not yet proven.
- Durable product persistence: projects/simulations/reports schema (migration
  0007) + tenant/owner-scoped ProductRepository + dual-write + scoped_project
  read-first + ArtifactStore abstraction; durable JobQueue + worker.py.
- SaaS hardening: durable RateLimiter (wired to login), UsageService (LLM
  accounting), redacted AuditService, idempotency, CORS allowlist, safe API
  errors, single-use PasswordResetService + endpoints (covers invite-pending).
- Exactly 3 roles (super_admin/admin/user) with tenant authz policy.
- Admin UI: GET/POST/PATCH /api/admin/users + GET/PUT /api/admin/settings
  (super-admin only, encrypted/masked); AdminView.vue + SettingsView.vue with
  admin/super-admin route guards, th/en i18n.
- Production deploy topology: multi-stage Dockerfile (frontend build + gunicorn
  wsgi + nginx SPA-proxy + supervisord worker), backend/wsgi.py, gunicorn dep.

Backend 197 passed; frontend 10 tests + build green. ruff unavailable (gap).
No commit of credentials; secrets handled via env/.env.example.
Deferred: Zep semantic A/B parity, object storage cutover, mobile QA, EasyPanel
container build of deploy topology.
2026-08-31 13:05:21 +07:00

817 lines
28 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
图谱相关API路由
采用项目上下文机制,服务端持久化状态
"""
import os
import threading
from flask import current_app, request, jsonify
from . import graph_bp
from ..config import Config
from ..services.ontology_generator import OntologyGenerator
from ..services.graph_builder import GraphBuilderService
from ..services.local_graph_builder import LocalGraphBuilderService
from ..utils.llm_client import LLMClient
from ..services.text_processor import TextProcessor
from ..utils.file_parser import FileParser
from ..utils.logger import get_logger
from ..utils.locale import t, get_locale, set_locale
from ..security.auth import current_actor, require_auth
from ..security.policy import Role
from ..utils.api_errors import ApiError
from ..models.task import TaskManager, TaskStatus
from ..models.project import ProjectManager, ProjectStatus
from ..services.product_repository import ProductRepository
from ..services.idempotency import idempotent
# 获取日志器
logger = get_logger('crowdsight.api')
@graph_bp.errorhandler(ApiError)
def handle_graph_error(error: ApiError):
return jsonify(error.to_payload(t)), error.status_code
def _scoped_project(project_id: str):
actor = current_actor()
owner_user_id = actor.user_id if actor.role is Role.USER else None
return ProjectManager.get_project_for_scope(
project_id,
organization_id=actor.organization_id,
owner_user_id=owner_user_id,
)
def _sync_project_to_durable(project, session_factory=None):
"""Best-effort dual-write of a project into the durable SQL repository.
Never raises: if no session factory is available (non-local backend, or a
background thread without one), the filesystem manager remains the
authoritative store and the durable copy is simply skipped. Callers inside a
worker thread pass the captured ``session_factory`` explicitly.
"""
try:
if session_factory is None:
session_factory = current_app.extensions.get("crowdsight_session_factory")
if session_factory is None:
return
session = session_factory()
try:
ProductRepository(session).sync_project(project, commit=True)
finally:
session.close()
except Exception:
logger.warning("durable project sync skipped", exc_info=True)
def _scoped_graph(graph_id: str):
actor = current_actor()
owner_user_id = actor.user_id if actor.role is Role.USER else None
return ProjectManager.find_project_by_graph_id(
graph_id,
organization_id=actor.organization_id,
owner_user_id=owner_user_id,
)
def _local_graph_builder(project):
"""Create the local graph adapter with the current tenant scope."""
session_factory = current_app.extensions.get("crowdsight_session_factory")
if not callable(session_factory):
raise ApiError("memory_backend_unavailable", 503, "api.internalError")
actor = current_actor()
return LocalGraphBuilderService(
session_factory,
organization_id=actor.organization_id,
project_id=project.project_id,
language=get_locale(),
)
def _scoped_task(task_id: str):
task = TaskManager().get_task(task_id)
if task is None:
return None
actor = current_actor()
metadata = task.metadata or {}
if metadata.get("organization_id") != actor.organization_id:
return None
if actor.role is Role.USER and metadata.get("owner_user_id") != actor.user_id:
return None
return task
def allowed_file(filename: str) -> bool:
"""ตรวจสอบนามสกุลไฟล์ที่อนุญาต"""
if not filename or '.' not in filename:
return False
ext = os.path.splitext(filename)[1].lower().lstrip('.')
return ext in Config.ALLOWED_EXTENSIONS
# ============== 项目管理接口 ==============
@graph_bp.route('/project/<project_id>', methods=['GET'])
@require_auth
def get_project(project_id: str):
"""
获取项目详情
"""
project = _scoped_project(project_id)
if not project:
return jsonify({
"success": False,
"error": t('api.projectNotFound', id=project_id)
}), 404
return jsonify({
"success": True,
"data": project.to_dict()
})
@graph_bp.route('/project/list', methods=['GET'])
@require_auth
def list_projects():
"""
列出所有项目
"""
limit = request.args.get('limit', 50, type=int)
actor = current_actor()
owner_user_id = actor.user_id if actor.role is Role.USER else None
projects = ProjectManager.list_projects(
limit=limit,
organization_id=actor.organization_id,
owner_user_id=owner_user_id,
)
return jsonify({
"success": True,
"data": [p.to_dict() for p in projects],
"count": len(projects)
})
@graph_bp.route('/project/<project_id>', methods=['DELETE'])
@require_auth
def delete_project(project_id: str):
"""
删除项目
"""
project = _scoped_project(project_id)
if project is None:
return jsonify({
"success": False,
"error": t('api.projectDeleteFailed', id=project_id)
}), 404
actor = current_actor()
success = ProjectManager.delete_project(
project.project_id,
organization_id=actor.organization_id,
owner_user_id=actor.user_id if actor.role is Role.USER else None,
)
if not success:
return jsonify({
"success": False,
"error": t('api.projectDeleteFailed', id=project_id)
}), 404
return jsonify({
"success": True,
"message": t('api.projectDeleted', id=project_id)
})
@graph_bp.route('/project/<project_id>/reset', methods=['POST'])
@require_auth
def reset_project(project_id: str):
"""
重置项目状态(用于重新构建图谱)
"""
project = _scoped_project(project_id)
if not project:
return jsonify({
"success": False,
"error": t('api.projectNotFound', id=project_id)
}), 404
# 重置到本体已生成状态
if project.ontology:
project.status = ProjectStatus.ONTOLOGY_GENERATED
else:
project.status = ProjectStatus.CREATED
project.graph_id = None
project.graph_build_task_id = None
project.error = None
ProjectManager.save_project(project)
return jsonify({
"success": True,
"message": t('api.projectReset', id=project_id),
"data": project.to_dict()
})
# ============== 接口1:上传文件并生成本体 ==============
@graph_bp.route('/ontology/generate', methods=['POST'])
@require_auth
@idempotent
def generate_ontology():
"""
接口1:上传文件,分析生成本体定义
请求方式:multipart/form-data
参数:
files: 上传的文件(PDF/MD/TXT),可多个
simulation_requirement: 模拟需求描述(必填)
project_name: 项目名称(可选)
additional_context: 额外说明(可选)
返回:
{
"success": true,
"data": {
"project_id": "proj_xxxx",
"ontology": {
"entity_types": [...],
"edge_types": [...],
"analysis_summary": "..."
},
"files": [...],
"total_text_length": 12345
}
}
"""
try:
logger.info("=== 开始生成本体定义 ===")
# 获取参数
simulation_requirement = request.form.get('simulation_requirement', '')
project_name = request.form.get('project_name', 'Unnamed Project')
additional_context = request.form.get('additional_context', '')
template_id = request.form.get('template_id', '')
logger.debug(f"项目名称: {project_name}")
logger.debug(f"模拟需求: {simulation_requirement[:100]}...")
if template_id:
logger.debug(f"Template: {template_id}")
# Load template filter rules if template_id provided
template_filter_rules = None
if template_id:
import json as _json
_templates_path = os.path.join(os.path.dirname(__file__), '..', 'templates.json')
try:
with open(_templates_path, 'r', encoding='utf-8') as _f:
_templates = _json.load(_f)['templates']
for _tmpl in _templates:
if _tmpl['id'] == template_id:
template_filter_rules = _tmpl.get('entity_filter', {})
break
except Exception:
pass
if not simulation_requirement:
return jsonify({
"success": False,
"error": t('api.requireSimulationRequirement')
}), 400
# 获取上传的文件
uploaded_files = request.files.getlist('files')
if not uploaded_files or all(not f.filename for f in uploaded_files):
return jsonify({
"success": False,
"error": t('api.requireFileUpload')
}), 400
# Create a tenant-owned project.
actor = current_actor()
project = ProjectManager.create_project(
name=project_name,
organization_id=actor.organization_id,
owner_user_id=actor.user_id,
)
project.simulation_requirement = simulation_requirement
# Best-effort dual-write into the durable repository.
_sync_project_to_durable(project)
logger.info(f"创建项目: {project.project_id}")
# 保存文件并提取文本
document_texts = []
all_text = ""
for file in uploaded_files:
if file and file.filename and allowed_file(file.filename):
# 保存文件到项目目录
file_info = ProjectManager.save_file_to_project(
project.project_id,
file,
file.filename
)
project.files.append({
"filename": file_info["original_filename"],
"size": file_info["size"]
})
# 提取文本
text = FileParser.extract_text(file_info["path"])
text = TextProcessor.preprocess_text(text)
document_texts.append(text)
all_text += f"\n\n=== {file_info['original_filename']} ===\n{text}"
if not document_texts:
ProjectManager.delete_project(
project.project_id,
organization_id=current_actor().organization_id,
owner_user_id=current_actor().user_id,
)
return jsonify({
"success": False,
"error": t('api.noDocProcessed')
}), 400
# 保存提取的文本
project.total_text_length = len(all_text)
ProjectManager.save_extracted_text(project.project_id, all_text)
logger.info(f"文本提取完成,共 {len(all_text)} 字符")
# 生成本体
logger.info("调用 LLM 生成本体定义...")
generator = OntologyGenerator()
ontology = generator.generate(
document_texts=document_texts,
simulation_requirement=simulation_requirement,
additional_context=additional_context if additional_context else None,
template_filter_rules=template_filter_rules
)
# 保存本体到项目
entity_count = len(ontology.get("entity_types", []))
edge_count = len(ontology.get("edge_types", []))
logger.info(f"本体生成完成: {entity_count} 个实体类型, {edge_count} 个关系类型")
project.ontology = {
"entity_types": ontology.get("entity_types", []),
"edge_types": ontology.get("edge_types", [])
}
project.analysis_summary = ontology.get("analysis_summary", "")
project.status = ProjectStatus.ONTOLOGY_GENERATED
ProjectManager.save_project(project)
_sync_project_to_durable(project)
logger.info(f"=== 本体生成完成 === 项目ID: {project.project_id}")
return jsonify({
"success": True,
"data": {
"project_id": project.project_id,
"project_name": project.name,
"ontology": project.ontology,
"analysis_summary": project.analysis_summary,
"files": project.files,
"total_text_length": project.total_text_length
}
})
except Exception as e:
logger.exception("Graph operation failed: %s", type(e).__name__)
raise ApiError("graph_operation_failed", 500, "api.internalError") from e
# ============== 接口2:构建图谱 ==============
@graph_bp.route('/build', methods=['POST'])
@require_auth
@idempotent
def build_graph():
"""
接口2:根据project_id构建图谱
请求(JSON):
{
"project_id": "proj_xxxx", // 必填,来自接口1
"graph_name": "图谱名称", // 可选
"chunk_size": 500, // 可选,默认500
"chunk_overlap": 50 // 可选,默认50
}
返回:
{
"success": true,
"data": {
"project_id": "proj_xxxx",
"task_id": "task_xxxx",
"message": "图谱构建任务已启动"
}
}
"""
try:
logger.info("=== 开始构建图谱 ===")
# Validate the selected backend before starting an asynchronous task.
backend = Config.MEMORY_BACKEND
if backend not in {"zep", "local"}:
raise ApiError("invalid_memory_backend", 500, "api.internalError")
errors = []
if backend == "zep" and not Config.ZEP_API_KEY:
errors.append(t('api.zepApiKeyMissing'))
if backend == "local" and not Config.LLM_API_KEY:
errors.append("LLM_API_KEY not configured")
if errors:
logger.error("Graph backend configuration is incomplete: backend=%s", backend)
return jsonify({
"success": False,
"error": t('api.configError', details="; ".join(errors))
}), 500
# 解析请求
data = request.get_json() or {}
project_id = data.get('project_id')
logger.debug(f"请求参数: project_id={project_id}")
if not project_id:
return jsonify({
"success": False,
"error": t('api.requireProjectId')
}), 400
# 获取项目
project = _scoped_project(project_id)
if not project:
return jsonify({
"success": False,
"error": t('api.projectNotFound', id=project_id)
}), 404
# 检查项目状态
force = data.get('force', False) # 强制重新构建
if project.status == ProjectStatus.CREATED:
return jsonify({
"success": False,
"error": t('api.ontologyNotGenerated')
}), 400
if project.status == ProjectStatus.GRAPH_BUILDING and not force:
return jsonify({
"success": False,
"error": t('api.graphBuilding'),
"task_id": project.graph_build_task_id
}), 400
# 如果强制重建,重置状态
if force and project.status in [ProjectStatus.GRAPH_BUILDING, ProjectStatus.FAILED, ProjectStatus.GRAPH_COMPLETED]:
project.status = ProjectStatus.ONTOLOGY_GENERATED
project.graph_id = None
project.graph_build_task_id = None
project.error = None
# 获取配置
graph_name = data.get('graph_name', project.name or 'CrowdSight Graph')
chunk_size = data.get('chunk_size', project.chunk_size or Config.DEFAULT_CHUNK_SIZE)
chunk_overlap = data.get('chunk_overlap', project.chunk_overlap or Config.DEFAULT_CHUNK_OVERLAP)
# 更新项目配置
project.chunk_size = chunk_size
project.chunk_overlap = chunk_overlap
# 获取提取的文本
text = ProjectManager.get_extracted_text(project_id)
if not text:
return jsonify({
"success": False,
"error": t('api.textNotFound')
}), 400
# 获取本体
ontology = project.ontology
if not ontology:
return jsonify({
"success": False,
"error": t('api.ontologyNotFound')
}), 400
# Create durable task metadata with tenant scope.
task_manager = TaskManager()
actor = current_actor()
organization_id = actor.organization_id
session_factory = current_app.extensions.get("crowdsight_session_factory")
if session_factory is None:
raise ApiError("memory_backend_unavailable", 503, "api.internalError")
task_id = task_manager.create_task(
f"构建图谱: {graph_name}",
metadata={
"project_id": project_id,
"organization_id": actor.organization_id,
"owner_user_id": actor.user_id,
},
)
logger.info(f"创建图谱构建任务: task_id={task_id}, project_id={project_id}")
# 更新项目状态
project.status = ProjectStatus.GRAPH_BUILDING
project.graph_build_task_id = task_id
ProjectManager.save_project(project)
_sync_project_to_durable(project)
# Capture locale before spawning background thread
current_locale = get_locale()
# 启动后台任务
def build_task():
set_locale(current_locale)
build_logger = get_logger('crowdsight.build')
try:
build_logger.info(f"[{task_id}] 开始构建图谱...")
task_manager.update_task(
task_id,
status=TaskStatus.PROCESSING,
message=t('progress.initGraphService')
)
# Choose the storage adapter inside the worker with captured scope.
if backend == "local":
builder = LocalGraphBuilderService(
session_factory,
organization_id=organization_id,
project_id=project_id,
extraction_client=LLMClient(),
language=current_locale,
)
else:
builder = GraphBuilderService(
api_key=Config.ZEP_API_KEY,
organization_id=organization_id,
owner_user_id=actor.user_id,
session_factory=session_factory,
)
# 分块
task_manager.update_task(
task_id,
message=t('progress.textChunking'),
progress=5
)
chunks = TextProcessor.split_text(
text,
chunk_size=chunk_size,
overlap=chunk_overlap
)
total_chunks = len(chunks)
# 创建图谱
task_manager.update_task(
task_id,
message=t('progress.creatingGraph' if backend == "local" else 'progress.creatingZepGraph'),
progress=10
)
graph_id = builder.create_graph(name=graph_name)
# 更新项目的graph_id
project.graph_id = graph_id
ProjectManager.save_project(project)
_sync_project_to_durable(project, session_factory)
# 设置本体
task_manager.update_task(
task_id,
message=t('progress.settingOntology'),
progress=15
)
builder.set_ontology(graph_id, ontology)
# 添加文本(progress_callback 签名是 (msg, progress_ratio))
def add_progress_callback(msg, progress_ratio):
progress = 15 + int(progress_ratio * 40) # 15% - 55%
task_manager.update_task(
task_id,
message=msg,
progress=progress
)
task_manager.update_task(
task_id,
message=t('progress.addingChunks', count=total_chunks),
progress=15
)
episode_uuids = builder.add_text_batches(
graph_id,
chunks,
batch_size=3,
progress_callback=add_progress_callback
)
# 等待Zep处理完成(查询每个episode的processed状态)
task_manager.update_task(
task_id,
message=t('progress.processingComplete' if backend == "local" else 'progress.waitingZepProcess'),
progress=55
)
def wait_progress_callback(msg, progress_ratio):
progress = 55 + int(progress_ratio * 35) # 55% - 90%
task_manager.update_task(
task_id,
message=msg,
progress=progress
)
builder._wait_for_episodes(episode_uuids, wait_progress_callback)
# 获取图谱数据
task_manager.update_task(
task_id,
message=t('progress.fetchingGraphData'),
progress=95
)
graph_data = builder.get_graph_data(graph_id)
# 更新项目状态
project.status = ProjectStatus.GRAPH_COMPLETED
ProjectManager.save_project(project)
_sync_project_to_durable(project, session_factory)
node_count = graph_data.get("node_count", 0)
edge_count = graph_data.get("edge_count", 0)
build_logger.info(f"[{task_id}] 图谱构建完成: graph_id={graph_id}, 节点={node_count}, 边={edge_count}")
# 完成
task_manager.update_task(
task_id,
status=TaskStatus.COMPLETED,
message=t('progress.graphBuildComplete'),
progress=100,
result={
"project_id": project_id,
"graph_id": graph_id,
"node_count": node_count,
"edge_count": edge_count,
"chunk_count": total_chunks
}
)
except Exception as e:
# 更新项目状态为失败
build_logger.error(
"[%s] graph build failed: error_type=%s",
task_id,
type(e).__name__,
)
project.status = ProjectStatus.FAILED
project.error = t('api.internalError')
ProjectManager.save_project(project)
_sync_project_to_durable(project, session_factory)
task_manager.update_task(
task_id,
status=TaskStatus.FAILED,
message=t('api.internalError'),
error=t('api.internalError'),
)
# 启动后台线程
thread = threading.Thread(target=build_task, daemon=True)
thread.start()
return jsonify({
"success": True,
"data": {
"project_id": project_id,
"task_id": task_id,
"message": t('api.graphBuildStarted', taskId=task_id)
}
})
except Exception as e:
logger.exception("Graph operation failed: %s", type(e).__name__)
raise ApiError("graph_operation_failed", 500, "api.internalError") from e
# ============== 任务查询接口 ==============
@graph_bp.route('/task/<task_id>', methods=['GET'])
@require_auth
def get_task(task_id: str):
"""
查询任务状态
"""
task = _scoped_task(task_id)
if not task:
return jsonify({
"success": False,
"error": t('api.taskNotFound', id=task_id)
}), 404
return jsonify({
"success": True,
"data": task.to_dict()
})
@graph_bp.route('/tasks', methods=['GET'])
@require_auth
def list_tasks():
"""
列出所有任务
"""
actor = current_actor()
tasks = [
task for task in TaskManager().list_tasks(
organization_id=actor.organization_id,
owner_user_id=actor.user_id if actor.role is Role.USER else None,
)
if _scoped_task(task.task_id) is not None
]
return jsonify({
"success": True,
"data": [t.to_dict() for t in tasks],
"count": len(tasks)
})
# ============== 图谱数据接口 ==============
@graph_bp.route('/data/<graph_id>', methods=['GET'])
@require_auth
def get_graph_data(graph_id: str):
"""
获取图谱数据(节点和边)
"""
project = _scoped_graph(graph_id)
if project is None:
return jsonify({
"success": False,
"error": t('api.graphNotBuilt'),
}), 404
try:
if Config.MEMORY_BACKEND == "local":
builder = _local_graph_builder(project)
elif Config.MEMORY_BACKEND == "zep":
if not Config.ZEP_API_KEY:
return jsonify({
"success": False,
"error": t('api.zepApiKeyMissing')
}), 500
builder = GraphBuilderService(api_key=Config.ZEP_API_KEY)
else:
raise ApiError("invalid_memory_backend", 500, "api.internalError")
graph_data = builder.get_graph_data(graph_id)
return jsonify({
"success": True,
"data": graph_data
})
except Exception as e:
logger.exception("Graph operation failed: %s", type(e).__name__)
raise ApiError("graph_operation_failed", 500, "api.internalError") from e
@graph_bp.route('/delete/<graph_id>', methods=['DELETE'])
@require_auth
def delete_graph(graph_id: str):
"""
ลบกราฟ memory
"""
project = _scoped_graph(graph_id)
if project is None:
return jsonify({
"success": False,
"error": t('api.graphNotBuilt'),
}), 404
try:
if Config.MEMORY_BACKEND == "local":
builder = _local_graph_builder(project)
elif Config.MEMORY_BACKEND == "zep":
if not Config.ZEP_API_KEY:
return jsonify({
"success": False,
"error": t('api.zepApiKeyMissing')
}), 500
builder = GraphBuilderService(api_key=Config.ZEP_API_KEY)
else:
raise ApiError("invalid_memory_backend", 500, "api.internalError")
builder.delete_graph(graph_id)
return jsonify({
"success": True,
"message": t('api.graphDeleted', id=graph_id)
})
except Exception as e:
logger.exception("Graph operation failed: %s", type(e).__name__)
raise ApiError("graph_operation_failed", 500, "api.internalError") from e