"""pipeline_llm.gateway — 模型治理引擎:门禁链 + 主备容错 + 账号端点轮转 + 双维度记账。 调用链(与设计规范 §2 一致): ① 个人限流窗口 → ② 组织限流窗口 → ③ 个人额度余额 → ④ 组织池余额(预授权冻结) → ⑤ 组织容错策略选模型(主→备) → ⑥ 账号×端点候选池:区域偏好过滤 → 余额加权轮转 + 限流冷却避让 → ⑦ 调用(调用方执行)→ ⑧ govern_settle 结算(实际用量,双维度记账) 集成方式(复用不改建): - llm_bridge._get_model_config 前置调 govern_resolve:机构配了策略/新模型表才走治理, 否则返回 None 由旧 llm 表逻辑兜底(向后兼容)。 - proxy_chat_completion(运行环境 token 路径)委托 chat_inference 统一结算。 - 结算唯一入口 govern_settle:usages 落上游 usage 原文 JSON(全量不截断, 计费唯一事实源;token 冗余列已删 m0021,2026-09-10 用户定夺)。 限流:Redis 分钟窗口原子计数(db4,与会话 db3 隔离)。Redis 不可用时限流放行(fail-open, 记 warning)——限流失效不应阻断业务,余额检查仍在 DB 侧守住。 余额:读-校验-条件 UPDATE(WHERE balance>=amount),并发竞态由 Redis 预授权计数托底。 错误必须真实可行动(用户铁律):每一级失败返回明确的中文原因,前端原样展示。 """ import json import logging import time logger = logging.getLogger("pipeline_llm.gateway") DBNAME = "pipeline" GOVERN_REDIS_DB = 4 # 治理计数器命名空间(会话用 db3,互不干扰) _COOLDOWN_SEC = 60 # 429 冷却窗口 _STICKY_TTL_SEC = 3600 # 会话粘性键 TTL(滑动续期,1h 无调用自动失效) _EST_TOKENS_PER_CALL = 800 # 历史常量(预授权已废弃,保留防外部引用) # 治理开关缓存:机构是否启用治理(有策略记录或有 llm_model 记录) _govern_cache: dict = {} _cfg_cache: dict = {} # (account_id) -> {api_key 解密, base_url...} # ────────────────────────── 基础设施 ────────────────────────── def _get_db(): from sqlor.dbpools import DBPools db = DBPools() if not db.databases: from appPublic.jsonConfig import getConfig config = getConfig() if config.databases: db.databases = config.databases return db, DBNAME def _password_key(): from appPublic.jsonConfig import getConfig try: return getConfig().password_key or 'QRIVSRHrthhwyjy176556332' except Exception: return 'QRIVSRHrthhwyjy176556332' def encrypt_api_key(plain: str) -> str: """AES 单层加密(对称可解密)。⚠️ 不用 RC4:password/unpassword 盐不对称,加密后不可解。""" if not plain: return '' from appPublic.aes import aes_encode_b64 return aes_encode_b64(_password_key(), plain) def decrypt_api_key(enc: str) -> str: """解密 api_key。非密文(历史明文/环境差异)原样返回。""" if not enc: return '' try: from appPublic.aes import aes_decode_b64 return aes_decode_b64(_password_key(), enc) except Exception: return enc def _redis(): """治理用 Redis 连接(db4)。从 session_redis 配置推导地址,不可用返回 None。""" try: import redis as _r from appPublic.jsonConfig import getConfig url = '' try: url = (getConfig().website or {}).get('session_redis', {}).get('url', '') except Exception: url = '' if not url: url = 'redis://127.0.0.1:6379/3' # 替换库号:会话用配置的库,治理固定 GOVERN_REDIS_DB base = url.rsplit('/', 1)[0] return _r.Redis.from_url(base + '/' + str(GOVERN_REDIS_DB), socket_timeout=2) except Exception as e: logger.warning("pipeline_llm: redis 不可用(限流降级放行): %s", e) return None def _row_to_dict(r): """行对象转 dict。sqlor 行是 DictObject(dict 子类),数据存在 dict 本体—— 用 vars() 取的是空的实例 __dict__(2026-09-04 实测根因),必须 dict() 拷贝。""" try: return dict(r) except (TypeError, ValueError): return {k: v for k, v in vars(r).items() if not callable(v)} async def _get_param(sor, name, default): """读 appbase params 表配置(params_name→params_value),表不存在/无值兜底 default。""" try: recs = await sor.sqlExe( "SELECT params_value FROM params WHERE params_name=${n}$ LIMIT 1", {"n": name}) if recs: v = str(getattr(recs[0], 'params_value', '') or '').strip() if v: return v except Exception: pass return default async def _wallet_threshold(sor) -> float: """账号钱包余额门禁阈值(2026-09-11 用户定夺:余额<阈值的账号从轮询集剔除 + 派待办给属主机构 operator 催充值)。params 可配,缺省 5 元。""" v = await _get_param(sor, 'llm_wallet_balance_threshold', '5') try: return float(v) except (TypeError, ValueError): return 5.0 def _parse_json(s, default): if not s: return default try: v = json.loads(s) return v if v is not None else default except Exception: return default def _fnum(v): try: return float(v or 0) except (TypeError, ValueError): return 0.0 # ────────────────────────── 门禁 ①②:限流 ────────────────────────── async def _rate_check(scope: str, scope_id: str, limit: int): """分钟窗口限流。返回 (ok, msg)。limit<=0 不限;Redis 故障放行。""" if not limit or limit <= 0 or not scope_id: return True, '' r = _redis() if r is None: return True, '' minute = int(time.time() // 60) key = 'llm_rl:%s:%s:%d' % (scope, scope_id, minute) try: cnt = r.incr(key) r.expire(key, 70) except Exception as e: logger.warning("pipeline_llm: 限流计数失败(放行): %s", e) return True, '' if cnt > limit: label = '个人' if scope == 'user' else '组织' return False, '%s限流:本分钟已调用 %d 次(上限 %d),请下一分钟再试' % (label, cnt, limit) return True, '' # ────────────────────────── 门禁 ⑤⑥:策略选模型 + 候选轮转 ────────────────────────── # 注:2026-09 记账重构废弃了组织池/个人额度预授权(原 _reserve_quota/_release_quota 已删除)—— # 本机构模型只记成本侧,owner 模型走异步三方账,防滥用由限流(①②)承担。 async def _org_policy(sor, org_id: str): """组织容错策略。无记录返回 (主模型空, 辅助模型空, 备链空, 端点偏好 any)。""" if not org_id: return '', '', [], 'any' recs = await sor.sqlExe( "SELECT primary_model_id, utility_model_id, backup_model_ids, endpoint_pref " "FROM llm_org_policy WHERE org_id=${o}$ AND status='active' LIMIT 1", {"o": org_id}) await sor.sqlExe("COMMIT", {}) if not recs: return '', '', [], 'any' r = recs[0] return (getattr(r, 'primary_model_id', '') or '', getattr(r, 'utility_model_id', '') or '', _parse_json(getattr(r, 'backup_model_ids', ''), []), getattr(r, 'endpoint_pref', '') or 'any') async def _model_chain(sor, org_id: str, model_name: str, purpose: str = ''): """解析模型链 [主, 备...],返回 llm_model 记录 dict 列表。 - 调用方指定 model_name:在 llm_model 找到 → [它] + **同能力**备链补充; 找不到 → 空(走旧表兜底) - 未指定:策略主模型 + 备链;无策略 → 空 - purpose='utility'(辅助任务:分类/选择/摘要):链 = [辅助模型] + 主 + 备。 显式指定 model_name 时忽略 purpose(显式优先);未配辅助模型时等同普通链。 辅助模型不可用(账号/端点全挂)由 govern_resolve 的逐模型尝试顺延到主模型, 与主备容错同一机制,无新增降级路径。 ⚠️ 同能力备链铁律(2026-09-07 qwen-image-plus 实测根因):显式指定 t2i 模型不可用时,旧逻辑把策略备链(清一色 t2t 对话模型)补进链里 → 静默 顺延到 deepseek 返回**文字**还报 ok:true——要图得文当成功,违反 「治理真实失败禁止静默回退」,且直接毒害 invoke_model 自动选型 (选了生图模型却拿到文本)。跨能力替换永远非法;同能力替代 (另一个文生图模型)才是合法容错。指定模型无同能力替代时链只剩 它自己,挑不到候选 → govern_resolve 报可行动错误(消息带最后原因)。 """ chain_ids = [] primary, utility, backups, pref = await _org_policy(sor, org_id) if model_name: recs = await sor.sqlExe( "SELECT id, capability FROM llm_model WHERE status='active' " "AND (name=${n}$ OR vendor_model_id=${n}$) LIMIT 1", {"n": model_name}) await sor.sqlExe("COMMIT", {}) if not recs: return [], pref chain_ids.append(getattr(recs[0], 'id', '')) want_cap = (getattr(recs[0], 'capability', '') or 't2t').strip().lower() if backups: # sqlor IN 列表必须展开占位符(传 list 会崩) ph = ','.join('${b%d}$' % i for i in range(len(backups))) bparams = {'b%d' % i: b for i, b in enumerate(backups)} # 能力降级(2026-09-08):备链允许同能力或合法降级源能力 # (t2i→i2i、t2v→r2v);跨能力(如 t2t)仍永远非法 from .selection import CAP_DEGRADE_SOURCE _alt = CAP_DEGRADE_SOURCE.get(want_cap, '') if _alt: ph += ',${altcap}$' bparams['altcap'] = _alt bparams['cap'] = want_cap _sql = ("SELECT id FROM llm_model WHERE status='active' AND id IN (%s) " "AND COALESCE(NULLIF(capability,''),'t2t') IN (${cap}$" + (",${altcap}$" if _alt else "") + ")") % ph brecs = await sor.sqlExe(_sql, bparams) await sor.sqlExe("COMMIT", {}) for br in (brecs or []): bid = getattr(br, 'id', '') if bid and bid not in chain_ids: chain_ids.append(bid) else: if purpose == 'utility' and utility: chain_ids.append(utility) if primary: chain_ids.append(primary) chain_ids.extend([b for b in backups if b not in chain_ids]) if not chain_ids: return [], pref models = [] for mid in chain_ids: recs = await sor.sqlExe( "SELECT id, name, vendor_id, vendor_model_id, capability, default_params, " "ppid, org_id, profile_id, sync_mode, query_profile_ids, account_id " "FROM llm_model WHERE id=${i}$ AND status='active' LIMIT 1", {"i": mid}) await sor.sqlExe("COMMIT", {}) if recs: models.append(_row_to_dict(recs[0])) return models, pref async def _vendor_endpoints(sor, vendor_id: str): """供应商端点目录。返回 [endpoint dict]。 2026-09-05:供应商表不再有 protocol 字段——请求形态(协议/模式)由模型挂的 适配模板(llm_api_profile)决定,供应商只提供端点目录(账号/区域/超时)。 """ recs = await sor.sqlExe( "SELECT endpoints FROM llm_vendor WHERE id=${i}$ AND status='active' LIMIT 1", {"i": vendor_id}) await sor.sqlExe("COMMIT", {}) if not recs: return [] return _parse_json(getattr(recs[0], 'endpoints', ''), []) async def _account_candidates(sor, model: dict): """模型的候选账号池:同供应商、启用中。 模型配置了 account_id 时(如某模型只允许用某个账号充值/密钥), 候选池收窄为该账号——模型级账号绑定优先于供应商级轮转。 """ where = "vendor_id=${v}$ AND status='active'" params = {"v": model.get('vendor_id', '')} bound = (model.get('account_id') or '').strip() if bound: where += " AND id=${a}$" params['a'] = bound recs = await sor.sqlExe( "SELECT id, name, api_key, endpoint_ids, balance, status FROM llm_account " "WHERE " + where, params) await sor.sqlExe("COMMIT", {}) return [_row_to_dict(r) for r in (recs or [])] def _filter_by_pref(candidates_with_ep, pref: str): """按端点偏好过滤/排序候选 (account, endpoint, idx) 列表。 must_*:只保留匹配区域,空则返回空(宁可冒泡不跨端点)。 prefer_*:匹配区域排前。any:原序。 """ if pref in ('must_domestic', 'prefer_domestic'): want = 'domestic' elif pref in ('must_intl', 'prefer_intl'): want = 'international' else: return candidates_with_ep matched = [c for c in candidates_with_ep if (c[1].get('region') or '') == want] if pref.startswith('must'): return matched rest = [c for c in candidates_with_ep if (c[1].get('region') or '') != want] return matched + rest async def _pick_candidate(sor, model: dict, pref: str, session_id: str = ''): """构建候选并选一个:协议过滤 → 区域偏好过滤 → 钱包余额门禁 → 冷却避让 → 会话粘性(同会话固定账号省上游缓存钱)→ 跨会话轮询。 账号钱包门禁(2026-09-11 用户定夺):llm_account.balance 是该账号的钱包 (充值入口 llm_account_recharge);出账成功后按 llm_usage.account_id 扣减。 balance < 阈值(params llm_wallet_balance_threshold,缺省 5 元)→ 账号从 轮询集剔除;余额状态派生待办给属主机构 operator 催充值(todos.py provider)。 全部剔除 → 报可行动错误(禁静默回退)。 会话粘性(2026-09-11 用户定夺):同一会话对同一模型固定同一账号——上游 KV/前缀缓存按账号(api_key)隔离,固定账号可命中缓存省钱。粘性键 llm_sticky:{session_id}:{model_id}(TTL 1h 滑动续期);粘性账号失效 (停用/冷却/余额不足)时自动重选并写回。无 session_id 的调用退化为纯轮询。 返回 (ok, result):ok 时 result=(account_dict, endpoint_dict);否则 result=错误信息。 """ vendor_id = model.get('vendor_id', '') if not vendor_id: return False, '模型「%s」未绑定供应商,无法选择账号' % model.get('name', '') endpoints = await _vendor_endpoints(sor, vendor_id) if not endpoints: return False, '供应商端点目录为空(模型治理→供应商→端点目录),请配置端点' accounts = await _account_candidates(sor, model) if not accounts: return False, '模型「%s」的供应商下无启用账号(或账号余额状态异常)' % model.get('name', '') # 端点 = 纯连接点(主机根+区域+超时),所有端点对所有模型可互换(2026-09-12 # 用户定夺:模型和端点不能强绑定,多端点=冗余/区域/多key,每个端点都支持 # 全部模型)。请求形态(协议/版本段/接口路径)完全由模型适配模板 profile # 的完整相对 path 承载:t2t → /compatible-mode/v1/chat/completions, # i2v → /api/v1/services/aigc/...(主机根 + 完整相对 path = 任意端点可拼)。 # 旧协议过滤已删(它正是间歇 404 根因:模型无 profile → 过滤跳过 → t2t 被 # 轮询到 api/v1 原生端点拼出 /chat/completions 404)。 # 展开 (账号, 端点) 对:账号选用的端点下标。 # endpoint_ids 空 = 选用全部端点(2026-09-12:端点可互换后,空不应再是 # 「跳过该账号」——历史数据里单端点供应商的账号常留空,语义即「用唯一端点」)。 pairs = [] for acc in accounts: idxs = _parse_json(acc.get('endpoint_ids'), []) if not idxs: idxs = list(range(len(endpoints))) for i in idxs: try: i = int(i) except (TypeError, ValueError): continue if 0 <= i < len(endpoints): pairs.append((acc, endpoints[i], i)) if not pairs: return False, '供应商端点目录与账号选用端点无交集——请检查模型治理→供应商端点目录与账号选用端点配置' pairs = _filter_by_pref(pairs, pref) if not pairs: return False, ('端点偏好为「仅 %s」,但没有可用的该区域端点。按合规要求不跨端点降级——' '请配置该区域端点或调整组织策略的端点偏好' % ('国内' if 'domestic' in pref else '国际')) # 账号钱包门禁(2026-09-11 用户定夺):余额 < 阈值的账号从轮询集剔除。 # 阈值 params 可配(llm_wallet_balance_threshold,缺省 5 元);余额不足由 # 待办 provider 派生待办催充值(同一事项只发一次=状态派生,充值到账自动消失)。 wallet_min = await _wallet_threshold(sor) low = [(acc, ep, i) for acc, ep, i in pairs if _fnum(acc.get('balance')) < wallet_min] pairs = [(acc, ep, i) for acc, ep, i in pairs if _fnum(acc.get('balance')) >= wallet_min] if not pairs: names = '、'.join(sorted({str(a.get('name') or a.get('id', '')) for a, _, _ in low})) return False, ('模型「%s」的全部候选账号钱包余额低于 %.2f 元(%s),已从轮询集剔除——' '请为账号充值(模型治理→供应商→供应商账号→充值)或调低门禁阈值' % (model.get('name', ''), wallet_min, names or '无')) r = _redis() def _in_cooldown(acc, i): if r is None: return False try: return bool(r.get('llm_cool:%s:%d' % (acc.get('id', ''), i)) or r.get('llm_cool:%s:*' % (acc.get('id', ''),))) except Exception: return False # 会话粘性(2026-09-11 用户定夺:同会话固定账号,省上游缓存钱): # 粘性账号仍在可用集且未冷却 → 直接用并续期;失效 → 落到轮询重选并写回。 sticky_key = '' if session_id: sticky_key = 'llm_sticky:%s:%s' % (session_id, model.get('id', '')) try: if r is not None: want = r.get(sticky_key) or '' if isinstance(want, bytes): want = want.decode() want = str(want).strip() if want: hit = [(acc, ep, i) for acc, ep, i in pairs if acc.get('id', '') == want and not _in_cooldown(acc, i)] if hit: r.set(sticky_key, want, ex=_STICKY_TTL_SEC) # 滑动续期 acc, ep, i = hit[0] return True, (acc, ep) except Exception as e: logger.warning("pipeline_llm: 会话粘性查询失败(退化为轮询): %s", e) # 冷却避让 + 跨会话轮询:Redis INCR 原子计数对可用集取模(同模型同供应商 # 各会话轮流用不同账号;Redis 不可用时退化为首个可用账号,fail-open)。 usable = [(acc, ep, i) for acc, ep, i in pairs if not _in_cooldown(acc, i)] if not usable: return False, ('模型「%s」的全部候选账号都在限流冷却中,请稍后重试' % model.get('name', '')) idx = 0 try: if r is not None: idx = int(r.incr('llm_rr:%s' % model.get('id', ''))) % len(usable) except Exception: idx = 0 acc, ep, i = usable[idx] if sticky_key: try: if r is not None: r.set(sticky_key, acc.get('id', ''), ex=_STICKY_TTL_SEC) except Exception: pass return True, (acc, ep) def _mark_cooldown(account_id: str, ep_index): """429 后把该账号(或账号×端点对)放进冷却。ep_index='*' 表示账号级。""" r = _redis() if r is None: return try: r.set('llm_cool:%s:%s' % (account_id, ep_index), '1', ex=_COOLDOWN_SEC) except Exception: pass def mark_cooldown_by_pair(account_id: str, base_url: str): """供代理层 429 回调:按 base_url 找到端点下标进冷却(找不到则全部下标冷却)。""" try: # 代理层只拿得到 base_url;冷却键按下标,这里无法反查下标—— # 退化为账号级冷却(该账号所有端点),宁可保守 _mark_cooldown(account_id, '*') except Exception: pass # ────────────────────────── 对外主入口 ────────────────────────── class GovernError(ValueError): """治理失败:消息真实可行动,前端原样展示。""" async def governance_enabled(sor, org_id: str) -> bool: """机构是否启用治理:**策略即开关**——该机构有 active 的容错策略才启用。 ⚠️ 不能用"新模型表有记录"判断(2026-09-01 修正):系统级(org_id=0)模型 注册后会让所有未配策略的机构被误判为启用,调用直接报错,断掉旧 llm 表链路。 机构显式建策略才是接入治理的动作;无策略 → 一切照旧(__LEGACY__)。 带缓存(改配置需重启,与 key 缓存一致)。 """ if not org_id: return False if org_id in _govern_cache: return _govern_cache[org_id] enabled = False try: recs = await sor.sqlExe( "SELECT COUNT(*) AS c FROM llm_org_policy WHERE org_id=${o}$ AND status='active'", {"o": org_id}) await sor.sqlExe("COMMIT", {}) enabled = (int(getattr(recs[0], 'c', 0)) if recs else 0) > 0 except Exception as e: logger.warning("pipeline_llm: governance_enabled 检查失败(按未启用): %s", e) enabled = False _govern_cache[org_id] = enabled return enabled async def _active_policy_org(sor, org_id: str) -> str: """确定生效策略机构:本机构有 active 策略 → 本机构; 否则平台 owner('0') 有 active 策略 → '0'(商业化默认:用平台资源即付费); 都没有 → ''(调用方回退旧 llm 表,零风险)。 """ if org_id and org_id != '0': recs = await sor.sqlExe( "SELECT COUNT(*) AS c FROM llm_org_policy WHERE org_id=${o}$ AND status='active'", {"o": org_id}) await sor.sqlExe("COMMIT", {}) if recs and int(_fnum(getattr(recs[0], 'c', 0))) > 0: return org_id recs = await sor.sqlExe( "SELECT COUNT(*) AS c FROM llm_org_policy WHERE org_id='0' AND status='active'", {}) await sor.sqlExe("COMMIT", {}) if recs and int(_fnum(getattr(recs[0], 'c', 0))) > 0: return '0' return '' async def govern_resolve(org_id: str, user_id: str = '', model_name: str = '', est_tokens: int = 0, task_ref: str = '', purpose: str = '', project_id: str = '', session_id: str = ''): """门禁链 ①-⑥。 策略来源(2026-09 商业化默认): 本机构有策略 → 用本机构模型链; 本机构无策略 → 用平台 owner('0') 策略(用平台资源即付费); owner 也没策略 → __LEGACY__ 旧表兜底。 本机构链路全不可用 → 容错顺延 owner 链(辅助→主→备同样机制)。 purpose='utility':辅助任务(分类/选择/摘要),链 = 辅助→主→备。显式 model_name 优先。 session_id(2026-09-11 用户定夺):会话粘性键——同一会话对同一模型固定同一 账号(上游 KV/前缀缓存按账号隔离,固定账号命中缓存省钱);跨会话轮询。 为空(如运行环境 token 路径只绑 project_id)退化为纯轮询,不报错。 记账决策(2026-09):**不做组织池预授权**—— 本机构模型只记成本侧(供应商账号扣减),不产生组织/个人消费; owner 模型三方账(客户/商户/供应商)走异步产品计费链路。 防滥用第一道闸保留:限流。 返回 (ok, result): ok=True → result = {api_base, api_key(明文,仅进程内), model_id, name, account_id, endpoint_region, reserved(恒0), model_row, user_id, org_id, task_ref, policy_org_id} ok=False → result = 错误信息(真实可行动) 无任何策略 → (False, '__LEGACY__'):调用方走旧 llm 表逻辑 """ db, dbname = _get_db() async with db.sqlorContext(dbname) as sor: policy_org = await _active_policy_org(sor, org_id or '') if not policy_org: return False, '__LEGACY__' # ① 个人限流 if user_id: recs = await sor.sqlExe( "SELECT rate_limit FROM llm_user_quota WHERE user_id=${u}$ AND org_id=${o}$ " "AND status='active' LIMIT 1", {"u": user_id, "o": org_id}) await sor.sqlExe("COMMIT", {}) ulimit = int(_fnum(getattr(recs[0], 'rate_limit', 0))) if recs else 0 ok, msg = await _rate_check('user', user_id, ulimit) if not ok: return False, msg # ② 组织限流 recs = await sor.sqlExe( "SELECT rate_limit FROM llm_org_quota WHERE org_id=${o}$ AND status='active' LIMIT 1", {"o": org_id}) await sor.sqlExe("COMMIT", {}) olimit = int(_fnum(getattr(recs[0], 'rate_limit', 0))) if recs else 0 ok, msg = await _rate_check('org', org_id, olimit) if not ok: return False, msg # ⑤ 模型链(含 ⑥ 偏好)—— 显式 model_name 时全局查找 models, pref = await _model_chain(sor, policy_org, model_name or '', purpose) if not models: if model_name: return False, '__LEGACY__' # 指定模型不在新表 → 旧表兜底 return False, ('生效策略未配置模型:请在模型治理→组织容错策略配置主/备模型,' '或在调用时指定模型名') # ⑤.5 模型归属硬校验(2026-09-07 用户需求):只有平台 owner('0'/系统级共享) # 或本机构的模型可调用;他机构模型从链中剔除,链内替代模型(策略备链)自然顺延。 # 显式指定的模型被剔除且链内无替代 → 报可行动错误(禁止静默回退/越权放行)。 from .selection import _owner_allowed allowed = [m for m in models if _owner_allowed(m.get('org_id'), org_id)] if len(allowed) != len(models): for m in models: if not _owner_allowed(m.get('org_id'), org_id): logger.info("pipeline_llm: 模型 %s 归属机构 %s,调用方 %s 无权使用," "已从模型链剔除(顺延替代模型)", m.get('name'), m.get('org_id'), org_id or '(未指定)') if not allowed: if model_name: return False, ('模型「%s」归属其他机构(非平台 owner 共享、也非本机构模型),' '不允许调用。请选择本机构或平台共享模型' '(模型治理→模型注册可查看归属)' % model_name) return False, ('生效策略模型链全部归属其他机构,调用方 %s 无可调用模型——' '请检查机构策略配置' % (org_id or '(未指定)')) models = allowed # ⑥ 逐模型尝试(辅助 → 主 → 备容错) last_err = '' chosen = None for m in models: ok, cand = await _pick_candidate(sor, m, pref, session_id) if ok and isinstance(cand, tuple): chosen = (m, cand) break last_err = cand if isinstance(cand, str) else last_err logger.info("pipeline_llm: 模型 %s 候选不可用(%s),尝试链内下一模型", m.get('name'), cand) # 本机构链全不可用 → 容错落 owner 链 if chosen is None and policy_org != '0': models0, pref0 = await _model_chain(sor, '0', model_name or '', purpose) # owner 容错链同样过归属校验(策略里可能配了他机构注册的模型) models0 = [m for m in models0 if _owner_allowed(m.get('org_id'), org_id)] # 同能力铁律(2026-09-07):显式指定模型时 owner 容错链也只许 # 同能力替代——要图给文的静默回退在任何一层都不合法。 # 能力降级(2026-09-08):合法降级源能力(t2i→i2i、t2v→r2v)除外。 if model_name and models0: from .selection import CAP_DEGRADE_SOURCE _want = (models0[0].get('capability') or 't2t').strip().lower() _alt = CAP_DEGRADE_SOURCE.get(_want, '') models0 = [m for m in models0 if (m.get('capability') or 't2t').strip().lower() in (_want, _alt)] for m in models0: ok, cand = await _pick_candidate(sor, m, pref0, session_id) if ok and isinstance(cand, tuple): chosen = (m, cand) policy_org = '0' pref = pref0 logger.info("pipeline_llm: 本机构模型链不可用,容错使用平台 owner 模型 %s", m.get('name')) break last_err = cand if isinstance(cand, str) else last_err if chosen is None: return False, '模型链全部不可用。最后原因:%s' % last_err model, chosen_pair = chosen acc, ep = chosen_pair # ⑦ 组装结果(无预授权,reserved 恒 0;记账在 govern_settle 按来源分流) return True, { 'api_base': (ep.get('base_url') or '').rstrip('/'), 'api_key': decrypt_api_key(acc.get('api_key') or ''), 'model_id': model.get('vendor_model_id') or model.get('name', ''), 'name': model.get('name', ''), 'account_id': acc.get('id', ''), 'endpoint_region': ep.get('region') or '', 'endpoint_proxy': ep.get('proxy') or '', 'timeout': int(_fnum(ep.get('timeout')) or 60), 'reserved': 0.0, 'model_row': model, 'user_id': user_id or '', 'org_id': org_id or '', 'task_ref': task_ref or '', 'project_id': project_id or '', 'session_id': session_id or '', 'policy_org_id': policy_org, } async def govern_settle(ctx: dict, ok_call: bool, note: str = '', usages: dict = None): """结算 ⑧:成本侧扣减 + 用量流水(三态记账状态机)。 ctx govern_resolve 成功时返回的 dict(含 policy_org_id) ok_call 上游调用是否成功(失败也记账:写 failed 流水) usages 上游返回的完整用量原文 JSON(落 llm_usage.usages,**全量不截断**, 2026-09-10 用户定夺):token 计数(prompt_tokens/completion_tokens/ prompt_tokens_details)与按量因子(视频时长/分辨率/图像张数)都在其中, 由异步出账读入定价引擎。**usages 是计费唯一事实源**——req_tokens/ resp_tokens 冗余列已删除(m0021),统计查询走 JSON_EXTRACT。 历史教训:[:2000] 截断产生非法 JSON → accounting json.loads 静默失败 → derived 求值兜底 0 → 永不计费,故禁止任何形式的截断。 accounting_status 三态(2026-09-05 用户定夺;2026-09-06 补 NULL 语义): created —— 已创建、待记账(仅 status='SUCCEEDED' 的行,出账循环拾取) accounted —— 记账成功 failed —— 记账失败 NULL —— 无记账意义:status≠SUCCEEDED 的行(FAILED/PENDING/RUNNING/recharge) 一律保持 NULL,不写 created(用户定夺:只有 SUCCEEDED 才有记账的意义) status 统一枚举(2026-09-06 用户定夺):SUCCEEDED/FAILED/PENDING/RUNNING(其他=recharge)。 自用 vs 跨机构的分流不在这里——product_accounting_generic 内部按 is_self_use 决定记账项(自用只记 PAY* 采购成本,跨机构另记 PAY 客户应付), 所有成功调用统一 created 进出账队列,成本账不遗漏。 """ if not ctx or not isinstance(ctx, dict): return db, dbname = _get_db() try: model = ctx.get('model_row') or {} # 2026-09-04:模型表不再存单价/成本四字段——成本与售价统一走 ppid 定价引擎, # 出账由异步记账循环(accounting.py)经产品层计算;此处 cost 恒 0。 # 2026-09-11 用户定夺:账号钱包(llm_account.balance)扣减在**产品记账 # 完成之后**——出账 worker 记账成功拿到 supplier_cost 再扣本流水 # account_id 对应账号的余额(accounting.py::_deduct_wallet);调用时点 # 不扣(真实成本只有定价引擎能算)。流水的 account_id 就是"调用模型的 # 钱包"指针,扣减与对账都按它定位。 cost = 0.0 caller_org = ctx.get('org_id') or '' async with db.sqlorContext(dbname) as sor: if not ok_call: # 上游限流(429)失败 → 账号级冷却,轮转避让下一轮不再选它 if '429' in str(note or ''): _mark_cooldown(ctx.get('account_id', ''), '*') from appPublic.uniqueID import getID usages_str = '' if usages: try: # 全量落库不截断(2026-09-10 用户定夺):截断产生非法 JSON → # accounting json.loads 静默失败 → derived 兜底 0 → 永不计费 usages_str = json.dumps(usages, ensure_ascii=False, default=str) except Exception: usages_str = '' # status 统一枚举(2026-09-06 用户定夺):SUCCEEDED/FAILED; # accounting_status 只在 SUCCEEDED 时写 created,否则缺省 NULL # (只有 SUCCEEDED 才有记账意义——失败/充值行挂 created 曾被误报"待记账") usage_row = { 'id': getID(), 'org_id': caller_org, 'user_id': ctx.get('user_id', ''), 'model_id': model.get('id', ''), 'account_id': ctx.get('account_id', ''), 'endpoint_region': ctx.get('endpoint_region', ''), 'cost': round(cost, 6), 'charge': 0, # 客户应付由异步出账计算(自用为 0,跨机构按折扣) 'ppid': model.get('ppid', '') or '', # 项目归属(2026-09-10 用户定夺):token 绑定的 project_id 透传, # 非项目调用(平台测试/充值/无项目 token)写哨兵'0'(列 NOT NULL # DEFAULT '0',查询不传=不区分项目),按项目统计费用 'project_id': ctx.get('project_id') or '0', 'task_ref': ctx.get('task_ref', ''), 'call_id': ctx.get('call_id', '') or '', 'status': 'SUCCEEDED' if ok_call else 'FAILED', 'usages': usages_str, 'note': (note or '')[:200], } if ok_call: usage_row['accounting_status'] = 'created' await sor.C('llm_usage', usage_row) await sor.sqlExe("COMMIT", {}) except Exception as e: logger.warning("pipeline_llm: 结算失败(不阻断调用方): %s", e) # ────────────────────────── 充值 ────────────────────────── async def recharge_account(account_id: str, amount, note: str = '', operator: str = ''): """供应商账号充值登记(成本侧钱包)。返回 (ok, msg)。 真实充值动作在供应商侧完成,本平台只做登记:balance += 金额、 total_recharge += 金额、写 recharge 流水。msg 带回变动前后值 (余额/累计充值 旧→新),前端无需再查即可展示「随之改变」。 """ try: amount = float(amount) except (TypeError, ValueError): return False, '充值金额必须是数字' if amount <= 0: return False, '充值金额必须大于 0' db, dbname = _get_db() async with db.sqlorContext(dbname) as sor: recs = await sor.sqlExe( "SELECT id, balance, total_recharge FROM llm_account WHERE id=${i}$", {"i": account_id}) await sor.sqlExe("COMMIT", {}) if not recs: return False, '账号不存在' old_bal = float(recs[0]['balance'] or 0) old_tot = float(recs[0]['total_recharge'] or 0) await sor.sqlExe( "UPDATE llm_account SET balance=balance+${a}$, total_recharge=total_recharge+${a}$ " "WHERE id=${i}$", {"a": amount, "i": account_id}) await sor.sqlExe("COMMIT", {}) await _write_recharge_usage(sor, '', '', account_id, amount, note, operator) _govern_cache.clear() return True, ('账号充值 %.4f 成功(余额 %.4f→%.4f,累计充值 %.4f→%.4f)' % (amount, old_bal, old_bal + amount, old_tot, old_tot + amount)) async def recharge_org(org_id: str, amount, note: str = '', operator: str = ''): """组织池充值(消费侧钱包)。无记录自动创建。返回 (ok, msg)。""" try: amount = float(amount) except (TypeError, ValueError): return False, '充值金额必须是数字' if amount <= 0: return False, '充值金额必须大于 0' if not org_id: return False, '缺少机构 ID' db, dbname = _get_db() async with db.sqlorContext(dbname) as sor: recs = await sor.sqlExe("SELECT id FROM llm_org_quota WHERE org_id=${o}$", {"o": org_id}) await sor.sqlExe("COMMIT", {}) if recs: await sor.sqlExe( "UPDATE llm_org_quota SET balance=balance+${a}$, total_recharge=total_recharge+${a}$ " "WHERE org_id=${o}$", {"a": amount, "o": org_id}) else: from appPublic.uniqueID import getID await sor.C('llm_org_quota', { 'id': getID(), 'org_id': org_id, 'balance': amount, 'total_recharge': amount, 'rate_limit': 0, 'status': 'active'}) await sor.sqlExe("COMMIT", {}) await _write_recharge_usage(sor, org_id, '', '', amount, note, operator) _govern_cache.clear() return True, '组织池充值 %.4f 成功' % amount async def _write_recharge_usage(sor, org_id, user_id, account_id, amount, note, operator): """充值流水(llm_usage 里 status=recharge 的行,对账用)。""" from appPublic.uniqueID import getID await sor.C('llm_usage', { 'id': getID(), 'org_id': org_id or '', 'user_id': user_id or '', 'model_id': '', 'account_id': account_id or '', 'endpoint_region': '', 'cost': round(float(amount), 6) if account_id else 0, 'charge': round(float(amount), 6) if org_id else 0, 'ppid': '', 'task_ref': 'recharge', 'status': 'recharge', 'note': ((note or '') + (' | 操作人:%s' % operator if operator else ''))[:200], }) await sor.sqlExe("COMMIT", {}) def load_gateway(): """注册治理函数到 ServerEnv(供 dspy/其他模块调用)。""" from ahserver.serverenv import ServerEnv env = ServerEnv() env.llm_govern_resolve = govern_resolve env.llm_govern_settle = govern_settle env.llm_recharge_account = recharge_account env.llm_recharge_org = recharge_org env.llm_encrypt_api_key = encrypt_api_key logger.info("[pipeline_llm] gateway loaded")