pipeline-service/pipeline_service/platform_model_tools.py
ymq ba65d586b1 feat(agent): 完备性硬门禁(2026-09-08用户定夺一C+二A)
意图识别归主循环LLM(system prompt语义引导),完备性硬门禁归代码。
- invoke_model入口check_completeness按能力契约校验必备输入:
  t2i/t2v/tts要描述,i2v/i2i要输入图,asr要音频,r2v要参考媒体
- 缺失返回QUESTION前缀,不调上游零费用(在llm_infer之前)
- v1角色agent+v2会话agent四处工具循环QUESTION硬拦截转
  ask_user/ask_question,禁回填靠LLM自觉(deepseek类会空参硬试烧钱)
- 22例单测全PASS(齐备/缺失/字符串vs数组形态/未登记能力不误伤)
2026-09-08 11:58:39 +08:00

299 lines
15 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.

"""pipeline_service.platform_model_tools — 平台模型工具(唯一实现,2026-09-07)。
用户需求:通用助手/产线 agent 可调用平台 owner 机构 + 本机构的全部
pipeline-llm 注册模型,按任务与用户输入自动匹配合适模型完成任务。
- 候选可见性/owner 硬校验/自动选型:pipeline_llm.selection
(models_catalog / auto_select_model,与推理治理链同一机构语义)
- 实际调用:llm_bridge.llm_infer → /pipeline-llm/api/v1/chat/completions
统一推理端点(门禁链 + 记账 + 同步/异步分流全在治理层,本层零旁路)
消费方(薄壳委托,禁止复制逻辑):
- agent_loop_v2.AgentExecutor._t_list_platform_models / _t_invoke_model
(会话 agent:驾驶舱 + 通用助手)
- agent_loop._exec_agent_tool(产线角色 agent v1)
"""
import json
import logging
logger = logging.getLogger("pipeline.platform_model_tools")
_CAP_DESC = {
"t2t": "文本对话", "i2t": "图像理解", "m2t": "多媒体理解",
"t2i": "文生图", "i2v": "图生视频", "t2v": "文生视频",
"r2v": "参考生视频", "tts": "语音合成", "asr": "语音识别",
"embedding": "向量化", "rerank": "重排序",
}
# ─────────────────────────────────────────────────────────────────
# 完备性硬门禁(2026-09-08 用户定夺:一C+二A)
#
# 用户原则:「意图识别后还要能识别用户输入是否有完成任务的完备条件,
# 如果不够,必须询问得到,不能做不好还做,浪费钱财」。
#
# 分工(一C):意图识别归主循环 LLM(system prompt 引导,零额外调用);
# **完备性硬门禁归代码**——放在贵工具入口、上游调用之前。理由:完备性是
# 能力契约的确定性检查(该能力要求哪些输入、传没传),不是语义猜测,
# 代码判定可靠且不依赖模型自觉(deepseek 类模型无视 prompt 指令有实测前科)。
#
# 范围(二A):只护花钱/不可逆动作(生成类调用真花钱)。闲聊/查询不设门禁。
#
# 返回 QUESTION: 前缀 → agent 循环把它当 ask_user 抛给用户(见
# agent_loop_v2 run loop / v1 的 ask 处理),不调上游、一分钱不花。
# ─────────────────────────────────────────────────────────────────
# 能力 → 必备输入契约(缺任一即不完备,必须先问用户要)
# 键名对应调用方 params 的三数组契约字段 / task 文本
_CAP_REQUIRED = {
"t2i": ("task",), # 文生图:必须有画面描述
"t2v": ("task",), # 文生视频:必须有画面描述
"tts": ("task",), # 语音合成:必须有要念的文本
"i2v": ("image_files",), # 图生视频:必须有输入图
"i2i": ("image_files",), # 图生图:必须有输入图
"asr": ("audio_files",), # 语音识别:必须有音频
"r2v": ("any_media",), # 参考生视频:至少一类参考媒体
}
_MEDIA_KEYS = ("image_files", "audio_files", "video_files")
# 必备输入的中文显示名(追问话术里给用户看,不用内部字段名)
_REQ_LABEL = {
"image_files": "输入图片",
"audio_files": "输入音频",
"video_files": "参考视频",
"any_media": "参考素材(图片/视频/音频任一)",
}
# task 类必备输入按能力的显示名(文生图要「画面描述」,语音合成要「朗读文本」)
_TASK_LABEL = {
"t2i": "画面描述", "t2v": "画面描述", "i2v": "运动/画面变化描述",
"tts": "朗读文本", "asr": "识别需求说明",
}
# 缺输入时对用户的追问话术(按能力,说清要什么、什么格式)
_CAP_ASK = {
"t2i": "生成图片需要先有画面描述。请说明:画什么主体、风格(写实/插画/商务示意图等)、"
"比例或尺寸、图中是否要有文字及内容。",
"t2v": "生成视频需要先有画面描述。请说明:画面内容、时长、分辨率、风格。",
"tts": "语音合成需要先有要朗读的文本。请提供文本内容,并说明音色/语速偏好(可选)。",
"i2v": "图生视频必须先提供输入图片(公网 URL 或上传)。另外请说明想要的运动/画面变化。",
"i2i": "图生图必须先提供输入图片(公网 URL 或上传)。另外请说明要怎么改(风格化/局部重绘等)。",
"asr": "语音识别必须先提供音频文件(公网 URL 或上传)。",
"r2v": "参考生视频必须先提供至少一类参考素材(图片/视频/音频,公网 URL 或上传),"
"并说明想要的画面内容与运动。",
}
def _biz_media(biz, key):
"""业务参数里的媒体数组是否非空(三数组契约,值可为字符串或数组)。"""
v = biz.get(key)
if isinstance(v, (list, tuple)):
return any(str(x or "").strip() for x in v)
return bool(str(v or "").strip())
def check_completeness(capability, task, biz):
"""完备性门禁:该能力契约要求的输入是否齐备。
返回 '' = 完备可执行;否则返回需要向用户追问的话术(QUESTION 内容)。
契约未登记的能力(如 embedding/新能力)不拦——宁可放行也别误伤,
但生成类主流能力全部登记在册。
"""
cap = (capability or "").strip().lower()
required = _CAP_REQUIRED.get(cap)
if not required:
return ""
missing = []
for req in required:
if req == "task":
if not (task or "").strip():
missing.append(_TASK_LABEL.get(cap, "画面/内容描述"))
elif req == "any_media":
if not any(_biz_media(biz or {}, k) for k in _MEDIA_KEYS):
missing.append(_REQ_LABEL["any_media"])
else:
if not _biz_media(biz or {}, req):
missing.append(_REQ_LABEL.get(req, req))
if missing:
ask = _CAP_ASK.get(cap, "该能力缺少必需输入,请补充:" + "、".join(missing))
return ("调用「" + _CAP_DESC.get(cap, cap) + "」模型缺少必需输入("
+ "、".join(missing) + "),已停在调用前、未产生任何费用。" + ask)
return ""
async def tool_list_platform_models(params, org_id):
"""列出平台可用模型(本机构+平台owner机构,含能力类型与描述)。"""
cap = str((params or {}).get("capability") or "").strip().lower()
try:
from pipeline_llm.selection import models_catalog
except ImportError:
return "FAIL: 模型治理模块(pipeline-llm)未安装,无法列出平台模型。"
try:
caps = (cap,) if cap else ()
models = await models_catalog(org_id or "0", capabilities=caps)
except Exception as e:
return "ERROR: 列模型失败: " + str(e)[:300]
if not models:
scope = ("能力 " + cap) if cap else "全部能力"
return ("平台当前无可用模型(" + scope + ";范围=本机构+平台owner机构)。"
"请在模型治理→模型注册中添加。")
lines = ["平台可用模型(" + str(len(models)) + " 个,本机构+平台owner机构):"]
for m in models:
cap_txt = m.get("capability") or "t2t"
cap_cn = _CAP_DESC.get(cap_txt, cap_txt)
desc = (":" + m["description"][:80]) if m.get("description") else ""
vendor = (" [" + m["vendor"] + "]") if m.get("vendor") else ""
lines.append(" - " + m.get("name", "") + vendor
+ "(" + cap_cn + "/" + cap_txt + ")" + desc)
return "\n".join(lines)
async def _resolve_capability(model_name, org_id):
"""按模型注册名反查能力类型(门禁需要知道该查哪份输入契约)。
走 models_catalog(与推理链同一机构可见性),查不到返回 ''。
"""
try:
from pipeline_llm.selection import models_catalog
models = await models_catalog(org_id or "0", capabilities=())
for m in models:
if model_name in (m.get("name"), m.get("vendor_model_id")):
return (m.get("capability") or "").strip().lower()
except Exception as e:
logger.warning("_resolve_capability(%s) 失败: %s", model_name, e)
return ""
async def tool_invoke_model(params, org_id, user_id=""):
"""调用平台模型完成生成类任务(文生图/视频/语音等非对话能力)。
流程:解析 task/model/capability/params → 未指定 model 时按 task 自动
选型(auto_select_model,候选=本机构+平台owner,capability 空=全能力)
→ **完备性硬门禁**(缺必需输入则 QUESTION 追问,不调上游不花钱)
→ 组包 payload → llm_infer 统一推理端点 → 提取生成物 URL 返回。
失败返回真实可行动错误(禁静默粉饰)。
"""
p = params or {}
task = str(p.get("task") or "").strip()
if not task:
return "FAIL: invoke_model 需要 task(任务描述/提示词)"
model = str(p.get("model") or "").strip()
cap = str(p.get("capability") or "").strip().lower()
# 业务参数(JSON 字符串 → dict)
biz = {}
raw_params = p.get("params")
if isinstance(raw_params, dict):
biz = dict(raw_params)
elif isinstance(raw_params, str) and raw_params.strip():
try:
parsed = json.loads(raw_params)
if isinstance(parsed, dict):
biz = parsed
except Exception:
return ("FAIL: params 不是合法 JSON 对象。媒体输入用三数组:"
"image_files/audio_files/video_files(URL或base64数组)")
# 自动选型:未显式指定 model 时按 task 匹配
if not model:
try:
from pipeline_llm.selection import auto_select_model
# capability 空 = 全能力候选,让匹配器按任务选最合适能力
model, reason = await auto_select_model(
org_id or "0", task, user_id=user_id or "",
capabilities=(cap if cap else ""))
if model:
logger.info("invoke_model 自动选型: %s(%s)org=%s",
model, reason, org_id or "0")
else:
return ("FAIL: 未能自动匹配到合适模型(" + str(reason) + ")。"
"可先用 list_platform_models 查看可用模型,"
"再用 model 参数指定。")
except ImportError:
return ("FAIL: 模型治理模块(pipeline-llm)未安装,"
"无法自动选型或调用平台模型。")
except Exception as e:
return "ERROR: 自动选型失败: " + str(e)[:300]
# ── 完备性硬门禁(2026-09-08 用户定夺:一C+二A)──
# 必须在 llm_infer 之前:缺输入直接 QUESTION 追问用户,上游一分钱不花。
# capability 未显式给出时按选中模型反查(门禁需知道查哪份输入契约)。
if not cap:
cap = await _resolve_capability(model, org_id)
gap = check_completeness(cap, task, biz)
if gap:
logger.info("invoke_model 完备性门禁拦截(%s): %s", cap or "?", gap[:80])
return "QUESTION: " + gap
# 组包:生成类模型从 messages 末条取 prompt 文本(inference._last_prompt_text)
payload = {"messages": [{"role": "user", "content": task}]}
payload.update(biz)
try:
from .llm_bridge import llm_infer
data = await llm_infer(
payload, model=model, org_id=org_id or "0",
user_id=user_id or "", timeout=600)
except Exception as e:
return "FAIL: 模型「" + model + "」调用失败:" + str(e)[:400]
# 提取生成物:media(本地持久 URL,优先)> choices[0].message.content
result_url = ""
media = data.get("media") if isinstance(data, dict) else None
if isinstance(media, dict):
result_url = (media.get("video") or media.get("image")
or media.get("audio") or media.get("glb")
or media.get("3dmodel") or "")
if not result_url and isinstance(data, dict):
try:
result_url = ((data.get("choices") or [{}])[0]
.get("message") or {}).get("content") or ""
except Exception:
result_url = ""
result_url = str(result_url or "").strip()
if result_url:
return ("OK: 模型「" + model + "」生成完成。产物地址:" + result_url)
# 无产物 URL:回原始输出摘要(真实,不粉饰)
try:
summary = json.dumps(data, ensure_ascii=False, default=str)[:500]
except Exception:
summary = str(data)[:500]
return ("模型「" + model + "」已调用但未返回可识别的产物地址。原始输出:"
+ summary)
# 平台模型工具的 v1 形态定义(AGENT_TOOLS 同款 {name,description,params}),
# 供 agent_loop.role_agent_run 追加到角色 agent 工具清单。
PLATFORM_MODEL_TOOLS_V1 = [
{
"name": "list_platform_models",
"description": ("列出平台当前可用的模型(本机构+平台owner机构的模型,含能力类型:"
"t2t对话/i2t图像理解/t2i文生图/t2v文生视频/i2v图生视频/tts语音合成/"
"asr语音识别等)。需要调用非对话能力(生图/视频/语音)前先查模型时用。"
"capability 参数可按能力过滤(如 t2i)"),
"params": {"capability": "可选:能力类型过滤(如 t2i/t2v/tts),空=全部"},
},
{
"name": "invoke_model",
"description": ("调用平台模型完成生成类任务(文生图/图生视频/文生视频/语音合成等"
"非对话能力;写代码/写文档等对话类工作由你自己完成,不要用本工具)。"
"model 可空——空时平台根据 task 自动匹配最合适的可用模型。"
"生成产物返回本地持久 URL,产物需要落盘时用 write_file 记录 URL"),
"params": {
"task": "任务描述/提示词(必填,如「一只在月球上弹吉他的猫」)",
"model": "可选:模型注册名(空=按任务自动匹配)",
"capability": "可选:能力类型(t2i/t2v/i2v/tts/asr等,自动匹配时用于过滤候选)",
"params": ("可选:业务参数 JSON 字符串。媒体输入用三数组契约:"
"image_files/audio_files/video_files(值为公网URL或base64的数组);"
"生成参数如 resolution/duration/size 按模型文档"),
},
},
]
async def exec_platform_model_tool(tool, params, org_id, user_id=""):
"""v1/v2 统一分发入口(薄壳)。"""
if tool == "list_platform_models":
return await tool_list_platform_models(params, org_id)
if tool == "invoke_model":
return await tool_invoke_model(params, org_id, user_id=user_id)
return "未实现: " + str(tool)