feat(llm-govern): llm_bridge/llm_proxy 挂治理钩子——门禁前置解析+成功/失败结算
- _resolve_cfg_with_govern: 机构有策略→门禁链(限流/限额/主备容错/端点选择/预授权); 无策略→__LEGACY__走旧llm表(向后兼容);治理真实失败抛错禁止静默回退 - llm_call/llm_call_msgs/llm_call_msgs_native 三处接入;失败结算释放预授权写failed流水 - proxy_chat_completion(运行环境token路径)同样接入,按上游真实usage结算
This commit is contained in:
parent
4fffee78c0
commit
6953a27937
@ -154,7 +154,81 @@ async def _post_chat_completion(url: str, headers: dict, payload: dict) -> dict:
|
||||
raise last_exc if last_exc else ValueError("LLM call failed")
|
||||
|
||||
|
||||
async def llm_call(prompt: str, model: str = None, temperature: float = 0.7, org_id: str = None) -> str:
|
||||
# ────────────────────────── pipeline_llm 治理钩子 ──────────────────────────
|
||||
# 机构配置了治理(llm_org_policy / llm_model 有记录)→ 调用走门禁链(限流/限额/
|
||||
# 主备容错/端点选择/预授权),调用后按实际用量结算(双维度记账)。
|
||||
# 未配置治理 → 返回 __LEGACY__ 走下方旧 llm 表逻辑(向后兼容,现有调用零改动)。
|
||||
# 治理真实失败(限流/余额不足/候选耗尽)→ 抛 ValueError,消息真实可行动,禁止静默回退
|
||||
# (否则治理被绕过,限额形同虚设)。
|
||||
|
||||
async def _resolve_cfg_with_govern(model, org_id, user_id, est_text):
|
||||
"""返回 (cfg, gctx)。gctx 非空 = 本次调用受治理(调用后须结算)。"""
|
||||
gctx = None
|
||||
try:
|
||||
from pipeline_llm.gateway import govern_resolve
|
||||
except ImportError:
|
||||
govern_resolve = None
|
||||
if govern_resolve is not None and org_id:
|
||||
try:
|
||||
est = max((len(est_text) if est_text else 0) // 2, 200)
|
||||
ok, res = await govern_resolve(
|
||||
org_id=org_id, user_id=user_id or '',
|
||||
model_name=model or '', est_tokens=est)
|
||||
if ok and isinstance(res, dict):
|
||||
gctx = res
|
||||
elif res != '__LEGACY__':
|
||||
raise ValueError(res)
|
||||
except ValueError:
|
||||
raise
|
||||
except Exception as e:
|
||||
logger.warning("llm_bridge: 治理前置异常(回退旧表): %s", e)
|
||||
if gctx:
|
||||
cfg = {"api_base": gctx["api_base"], "api_key": gctx["api_key"],
|
||||
"model_id": gctx["model_id"]}
|
||||
return cfg, gctx
|
||||
cfg = await _get_model_config(model, org_id=org_id)
|
||||
return cfg, None
|
||||
|
||||
|
||||
async def _settle_govern(gctx, data, est_text):
|
||||
"""调用成功后结算:优先上游真实 usage,缺失则按文本长度估算(note=est)。"""
|
||||
if not gctx:
|
||||
return
|
||||
try:
|
||||
from pipeline_llm.gateway import govern_settle
|
||||
usage = (data or {}).get('usage') or {}
|
||||
rt = int(usage.get('prompt_tokens') or 0)
|
||||
ct = int(usage.get('completion_tokens') or 0)
|
||||
note = ''
|
||||
if not usage:
|
||||
rt = rt or max((len(est_text) if est_text else 0) // 2, 100)
|
||||
ct = ct or 200
|
||||
note = 'est'
|
||||
await govern_settle(gctx, True, rt, ct, note)
|
||||
except Exception as e:
|
||||
logger.warning("llm_bridge: 治理结算失败(不阻断调用): %s", e)
|
||||
|
||||
|
||||
async def _settle_govern_failed(gctx, note):
|
||||
"""调用失败时结算:释放预授权,写 failed 流水。"""
|
||||
if not gctx:
|
||||
return
|
||||
try:
|
||||
from pipeline_llm.gateway import govern_settle
|
||||
await govern_settle(gctx, False, 0, 0, str(note)[:200])
|
||||
except Exception as e:
|
||||
logger.warning("llm_bridge: 治理失败结算异常: %s", e)
|
||||
|
||||
|
||||
def _msgs_text(messages):
|
||||
try:
|
||||
return ' '.join(str(m.get('content', '')) for m in (messages or []) if isinstance(m, dict))
|
||||
except Exception:
|
||||
return ''
|
||||
|
||||
|
||||
async def llm_call(prompt: str, model: str = None, temperature: float = 0.7,
|
||||
org_id: str = None, user_id: str = None) -> str:
|
||||
"""Call LLM and return text response.
|
||||
|
||||
Backend priority:
|
||||
@ -176,13 +250,14 @@ async def llm_call(prompt: str, model: str = None, temperature: float = 0.7, org
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# Priority 2: DB llm table
|
||||
cfg = await _get_model_config(model, org_id=org_id)
|
||||
# Priority 2: DB llm table(治理启用时前置走门禁链)
|
||||
cfg, gctx = await _resolve_cfg_with_govern(model, org_id, user_id, prompt)
|
||||
if cfg.get("api_key") and cfg.get("api_base"):
|
||||
api_base = cfg["api_base"]
|
||||
api_key = cfg["api_key"]
|
||||
model_id = cfg.get("model_id") or model or "default"
|
||||
logger.info("llm_bridge: using DB model config for %s -> %s", model, api_base)
|
||||
logger.info("llm_bridge: using %s model config for %s -> %s",
|
||||
"governed" if gctx else "DB", model, api_base)
|
||||
else:
|
||||
# Priority 3: Environment variables
|
||||
api_base = os.environ.get("LLM_API_BASE", "https://api.openai.com/v1")
|
||||
@ -205,7 +280,12 @@ async def llm_call(prompt: str, model: str = None, temperature: float = 0.7, org
|
||||
}
|
||||
|
||||
url = api_base.rstrip("/") + "/chat/completions"
|
||||
data = await _post_chat_completion(url, headers, payload)
|
||||
try:
|
||||
data = await _post_chat_completion(url, headers, payload)
|
||||
except Exception as e:
|
||||
await _settle_govern_failed(gctx, e)
|
||||
raise
|
||||
await _settle_govern(gctx, data, prompt)
|
||||
return data["choices"][0]["message"]["content"]
|
||||
|
||||
|
||||
@ -214,11 +294,12 @@ async def call_llm(tenant_id: str, prompt: str, model: str = None, temperature:
|
||||
return await llm_call(prompt, model=model, temperature=temperature)
|
||||
|
||||
|
||||
async def llm_call_msgs(messages: list, model: str = None, temperature: float = 0.7, org_id: str = None) -> str:
|
||||
async def llm_call_msgs(messages: list, model: str = None, temperature: float = 0.7,
|
||||
org_id: str = None, user_id: str = None) -> str:
|
||||
"""Call LLM with full message array (system/user/assistant)."""
|
||||
import aiohttp
|
||||
|
||||
cfg = await _get_model_config(model, org_id=org_id)
|
||||
cfg, gctx = await _resolve_cfg_with_govern(model, org_id, user_id, _msgs_text(messages))
|
||||
if cfg.get("api_key") and cfg.get("api_base"):
|
||||
api_base = cfg["api_base"]
|
||||
api_key = cfg["api_key"]
|
||||
@ -235,11 +316,18 @@ async def llm_call_msgs(messages: list, model: str = None, temperature: float =
|
||||
payload = {"model": model_id, "messages": messages, "temperature": temperature}
|
||||
|
||||
url = api_base.rstrip("/") + "/chat/completions"
|
||||
data = await _post_chat_completion(url, headers, payload)
|
||||
try:
|
||||
data = await _post_chat_completion(url, headers, payload)
|
||||
except Exception as e:
|
||||
await _settle_govern_failed(gctx, e)
|
||||
raise
|
||||
await _settle_govern(gctx, data, _msgs_text(messages))
|
||||
return data["choices"][0]["message"]["content"]
|
||||
|
||||
|
||||
async def llm_call_msgs_native(messages: list, tools: list = None, model: str = None, temperature: float = 0.7, org_id: str = None) -> dict:
|
||||
async def llm_call_msgs_native(messages: list, tools: list = None, model: str = None,
|
||||
temperature: float = 0.7, org_id: str = None,
|
||||
user_id: str = None) -> dict:
|
||||
"""Native function calling. 传入 tools JSON schema,返回 message dict。
|
||||
|
||||
Returns:
|
||||
@ -248,7 +336,7 @@ async def llm_call_msgs_native(messages: list, tools: list = None, model: str =
|
||||
"""
|
||||
import aiohttp
|
||||
|
||||
cfg = await _get_model_config(model, org_id=org_id)
|
||||
cfg, gctx = await _resolve_cfg_with_govern(model, org_id, user_id, _msgs_text(messages))
|
||||
if cfg.get("api_key") and cfg.get("api_base"):
|
||||
api_base = cfg["api_base"]
|
||||
api_key = cfg["api_key"]
|
||||
@ -268,7 +356,12 @@ async def llm_call_msgs_native(messages: list, tools: list = None, model: str =
|
||||
payload["tool_choice"] = "auto"
|
||||
|
||||
url = api_base.rstrip("/") + "/chat/completions"
|
||||
data = await _post_chat_completion(url, headers, payload)
|
||||
try:
|
||||
data = await _post_chat_completion(url, headers, payload)
|
||||
except Exception as e:
|
||||
await _settle_govern_failed(gctx, e)
|
||||
raise
|
||||
await _settle_govern(gctx, data, _msgs_text(messages))
|
||||
msg = data["choices"][0]["message"]
|
||||
return {
|
||||
"content": msg.get("content") or "",
|
||||
|
||||
@ -275,17 +275,45 @@ async def proxy_chat_completion(token, payload):
|
||||
# 否则用请求里的 model;都没有则由 llm_bridge 取该机构第一个 active 模型。
|
||||
model_name = info.get('model_name') or payload.get('model') or None
|
||||
|
||||
from .llm_bridge import _get_model_config, _post_chat_completion
|
||||
cfg = await _get_model_config(model_name, org_id=org_id)
|
||||
if not (cfg.get('api_key') and cfg.get('api_base')):
|
||||
return False, f"机构 {org_id} 未配置可用模型(llm 表 status=active)"
|
||||
# pipeline_llm 治理前置:机构配了治理 → 门禁链(限流/限额/主备容错/端点选择/预授权);
|
||||
# 未配置 → __LEGACY__ 走旧 llm 表。治理真实失败(限流/余额不足)→ 抛错,禁止静默回退。
|
||||
gctx = None
|
||||
try:
|
||||
from pipeline_llm.gateway import govern_resolve
|
||||
_est = 200
|
||||
try:
|
||||
_msgs = payload.get('messages') or []
|
||||
_est = max(sum(len(str(m.get('content', ''))) for m in _msgs if isinstance(m, dict)) // 2, 200)
|
||||
except Exception:
|
||||
_est = 200
|
||||
gok, gres = await govern_resolve(
|
||||
org_id=org_id, user_id='', model_name=model_name or '', est_tokens=_est,
|
||||
task_ref='proxy:%s' % (info.get('project_id') or ''))
|
||||
if gok and isinstance(gres, dict):
|
||||
gctx = gres
|
||||
elif gres != '__LEGACY__':
|
||||
return False, gres
|
||||
except ImportError:
|
||||
gctx = None
|
||||
except Exception as e:
|
||||
logger.warning("proxy_chat_completion: 治理前置异常(回退旧表): %s", e)
|
||||
|
||||
api_base = cfg['api_base']
|
||||
api_key = cfg['api_key'] # 只在本进程内存中使用,不返回给调用方
|
||||
model_id = cfg.get('model_id') or model_name or 'default'
|
||||
if gctx:
|
||||
api_base = gctx['api_base']
|
||||
api_key = gctx['api_key']
|
||||
model_id = gctx['model_id']
|
||||
else:
|
||||
from .llm_bridge import _get_model_config
|
||||
cfg = await _get_model_config(model_name, org_id=org_id)
|
||||
if not (cfg.get('api_key') and cfg.get('api_base')):
|
||||
return False, f"机构 {org_id} 未配置可用模型(模型治理/llm 表 status=active)"
|
||||
api_base = cfg['api_base']
|
||||
api_key = cfg['api_key'] # 只在本进程内存中使用,不返回给调用方
|
||||
model_id = cfg.get('model_id') or model_name or 'default'
|
||||
|
||||
# 透传客户端参数(messages/tools/temperature/max_tokens 等),但 model 换成真实 model_id,
|
||||
# 且不透传 stream(代理暂不支持流式)。
|
||||
from .llm_bridge import _post_chat_completion
|
||||
upstream = {k: v for k, v in payload.items() if k not in ('model', 'stream')}
|
||||
upstream['model'] = model_id
|
||||
|
||||
@ -295,11 +323,31 @@ async def proxy_chat_completion(token, payload):
|
||||
data = await _post_chat_completion(url, headers, upstream)
|
||||
except Exception as e:
|
||||
logger.warning("proxy_chat_completion upstream failed: %s", e)
|
||||
# 治理路径:失败结算(释放预授权,写 failed 流水)
|
||||
if gctx:
|
||||
try:
|
||||
from pipeline_llm.gateway import govern_settle
|
||||
await govern_settle(gctx, False, 0, 0, str(e)[:200])
|
||||
except Exception:
|
||||
pass
|
||||
return False, f"上游调用失败: {str(e)[:300]}"
|
||||
|
||||
await _record_usage(info['id'], data.get('usage'))
|
||||
logger.info("proxy_chat_completion: org=%s project=%s model=%s ok",
|
||||
org_id, info.get('project_id'), model_id)
|
||||
# 治理路径:按上游真实 usage 结算(双维度记账)
|
||||
if gctx:
|
||||
try:
|
||||
from pipeline_llm.gateway import govern_settle
|
||||
_u = data.get('usage') or {}
|
||||
await govern_settle(
|
||||
gctx, True,
|
||||
int(_u.get('prompt_tokens') or 0),
|
||||
int(_u.get('completion_tokens') or 0),
|
||||
'' if _u else 'est')
|
||||
except Exception as e:
|
||||
logger.warning("proxy_chat_completion: 治理结算失败(不阻断): %s", e)
|
||||
logger.info("proxy_chat_completion: org=%s project=%s model=%s ok%s",
|
||||
org_id, info.get('project_id'), model_id,
|
||||
" (governed)" if gctx else "")
|
||||
return True, data
|
||||
|
||||
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user