diff --git a/pipeline_service/llm_bridge.py b/pipeline_service/llm_bridge.py index 53ed366..ef95e33 100644 --- a/pipeline_service/llm_bridge.py +++ b/pipeline_service/llm_bridge.py @@ -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 "", diff --git a/pipeline_service/llm_proxy.py b/pipeline_service/llm_proxy.py index cc1c2a0..eed2373 100644 --- a/pipeline_service/llm_proxy.py +++ b/pipeline_service/llm_proxy.py @@ -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