fix: use standalone get_sor_context(env, modulename)

This commit is contained in:
yumoqing 2026-07-22 15:25:42 +08:00
parent 01a2c9281f
commit abba9d1273

79
init.py
View File

@ -6,12 +6,12 @@ from traceback import format_exc
from ahserver.serverenv import ServerEnv from ahserver.serverenv import ServerEnv
from appPublic.registerfunction import RegisterFunction from appPublic.registerfunction import RegisterFunction
from appPublic.log import debug, exception from appPublic.log import debug, exception
from sqlor.dbpools import get_sor_context
import json import json
async def status_handler(request, params_kw, *args, **kwargs): async def status_handler(request, params_kw, *args, **kwargs):
"""服务状态""" """服务状态"""
env = request._run_ns
return json.dumps({ return json.dumps({
"service": "ragserver", "service": "ragserver",
"version": "0.1.0", "version": "0.1.0",
@ -27,12 +27,12 @@ async def kb_list_handler(request, params_kw, *args, **kwargs):
env = request._run_ns env = request._run_ns
try: try:
userorgid = await env.get_userorgid() userorgid = await env.get_userorgid()
sor = await env.get_sor_context('rag') async with get_sor_context(env, 'rag') as sor:
recs = await sor.R("knowledge_bases", {"org_id": userorgid}) recs = await sor.R("knowledge_bases", {"org_id": userorgid})
return json.dumps({ return json.dumps({
"status": "SUCCEEDED", "status": "SUCCEEDED",
"knowledge_bases": [dict(r) for r in recs] "knowledge_bases": [dict(r) for r in recs]
}, ensure_ascii=False, default=str) }, ensure_ascii=False, default=str)
except Exception as e: except Exception as e:
exception(f"kb_list: {e}, {format_exc()}") exception(f"kb_list: {e}, {format_exc()}")
return json.dumps({"error": str(e)}) return json.dumps({"error": str(e)})
@ -43,18 +43,17 @@ async def engines_handler(request, params_kw, *args, **kwargs):
env = request._run_ns env = request._run_ns
try: try:
userorgid = await env.get_userorgid() userorgid = await env.get_userorgid()
sor = await env.get_sor_context('rag') async with get_sor_context(env, 'rag') as sor:
# 返回全局默认 + 租户自定义,租户优先 sql = """
sql = """ SELECT * FROM engine_configs
SELECT * FROM engine_configs WHERE status='active' AND (org_id IS NULL OR org_id=%s)
WHERE status='active' AND (org_id IS NULL OR org_id=%s) ORDER BY engine_type, priority DESC
ORDER BY engine_type, priority DESC """
""" recs = await sor.sqlExe(sql, {"org_id": userorgid})
recs = await sor.sqlExe(sql, {"org_id": userorgid}) return json.dumps({
return json.dumps({ "status": "SUCCEEDED",
"status": "SUCCEEDED", "engines": [dict(r) for r in recs]
"engines": [dict(r) for r in recs] }, ensure_ascii=False, default=str)
}, ensure_ascii=False, default=str)
except Exception as e: except Exception as e:
exception(f"engines: {e}, {format_exc()}") exception(f"engines: {e}, {format_exc()}")
return json.dumps({"error": str(e)}) return json.dumps({"error": str(e)})
@ -67,60 +66,38 @@ async def search_handler(request, params_kw, *args, **kwargs):
query = params_kw.get("query", "") query = params_kw.get("query", "")
kb_id = params_kw.get("kb_id", "") kb_id = params_kw.get("kb_id", "")
top_k = int(params_kw.get("top_k", 5)) top_k = int(params_kw.get("top_k", 5))
if not query: if not query:
return json.dumps({"error": "query required"}) return json.dumps({"error": "query required"})
# 尝试调用 rag-pipeline 引擎
try: try:
from pipeline import search as pipeline_search from pipeline import search as pipeline_search
result = pipeline_search( result = pipeline_search(query, pipeline_name="kg-rag-standard",
query, collection=kb_id or "knowledge",
pipeline_name=params_kw.get("pipeline", "kg-rag-standard"), graph_name=kb_id or "knowledge",
collection=kb_id or "knowledge", top_k=top_k, llm_func=None)
graph_name=kb_id or "knowledge",
top_k=top_k,
llm_func=None
)
return json.dumps(result, ensure_ascii=False) return json.dumps(result, ensure_ascii=False)
except ImportError: except ImportError:
return json.dumps({ return json.dumps({"status": "FALLBACK", "message": "pipeline not available"})
"status": "FALLBACK",
"message": "rag-pipeline not available, search via direct API calls",
"query": query
}, ensure_ascii=False)
except Exception as e: except Exception as e:
exception(f"search: {e}, {format_exc()}") exception(f"search: {e}, {format_exc()}")
return json.dumps({"error": str(e)}) return json.dumps({"error": str(e)})
async def ingest_handler(request, params_kw, *args, **kwargs): async def ingest_handler(request, params_kw, *args, **kwargs):
"""文档入库 (文本直接调用 pipeline)""" """文档入库"""
env = request._run_ns env = request._run_ns
try: try:
document = params_kw.get("document", "") document = params_kw.get("document", "")
kb_id = params_kw.get("kb_id", "") kb_id = params_kw.get("kb_id", "")
if not document: if not document:
return json.dumps({"error": "document text required"}) return json.dumps({"error": "document text required"})
try: try:
from pipeline import ingest as pipeline_ingest from pipeline import ingest as pipeline_ingest
result = pipeline_ingest( result = pipeline_ingest(document, pipeline_name="kg-rag-standard",
document, collection=kb_id or "knowledge",
pipeline_name=params_kw.get("pipeline", "kg-rag-standard"), graph_name=kb_id or "knowledge", llm_func=None)
collection=kb_id or "knowledge",
graph_name=kb_id or "knowledge",
llm_func=None
)
return json.dumps(result, ensure_ascii=False) return json.dumps(result, ensure_ascii=False)
except ImportError: except ImportError:
return json.dumps({ return json.dumps({"status": "FALLBACK", "message": "pipeline not available"})
"status": "FALLBACK",
"message": "rag-pipeline not available for ingest"
}, ensure_ascii=False)
except Exception as e: except Exception as e:
exception(f"ingest: {e}, {format_exc()}") exception(f"ingest: {e}, {format_exc()}")
return json.dumps({"error": str(e)}) return json.dumps({"error": str(e)})