feat(bridge): session_id会话粘性透传(2026-09-11用户定夺:同会话固定账号省上游缓存钱,跨会话轮询)——llm_bridge四函数(llm_call/llm_call_msgs/llm_call_msgs_native/llm_infer)加session_id参数经payload._session_id透传(不进token缓存键);agent_loop_v2七处调用点接self.session_id(主循环native+降级重试+文本回退/技能选择/项目匹配/历史摘要/invoke_model);agent_loop v1四处任务循环用task:{task_id}做粘性键(任务内30轮同前缀);platform_model_tools.tool_invoke_model透传

This commit is contained in:
ymq 2026-09-11 15:30:30 +08:00
parent 58853de72d
commit 5c063ca588
4 changed files with 66 additions and 30 deletions

View File

@ -2034,7 +2034,8 @@ async def _exec_agent_tool(tool, params, workspace_dir, ctx=None):
from .platform_model_tools import exec_platform_model_tool
return await exec_platform_model_tool(
tool, p, org_id=str((ctx or {}).get('org_id', '') or '0'),
user_id=str((ctx or {}).get('user_id', '') or ''))
user_id=str((ctx or {}).get('user_id', '') or ''),
project_id=str((ctx or {}).get('project_id', '') or ''))
if tool == 'read_file':
path = p.get('path', '')
if not path: return 'FAIL: 需要文件路径'
@ -2475,7 +2476,7 @@ async def role_agent_run(project_id, role, agent_id=None, model_name=None):
msgs.append({"role": "user", "content": _FORCE_PRODUCE_HINT})
try:
resp = await asyncio.wait_for(
llm_call_msgs_native(msgs, tools=tools_schema, model=model_name, temperature=0.4, org_id=org_id),
llm_call_msgs_native(msgs, tools=tools_schema, model=model_name, temperature=0.4, org_id=org_id, project_id=project_id, session_id='task:%s' % task_id),
timeout=_LLM_HARD_TIMEOUT)
except Exception as e:
err_msg = f"{type(e).__name__}: {str(e)[:400]}"
@ -3061,7 +3062,7 @@ async def pm_review_run(project_id, agent_id=None, model_name=None):
"随后必须立即输出 review_approve / review_reject / review_complete / review_rollback 之一,"
"禁止再调用其它工具。"})
try:
raw = await llm_call_msgs(msgs, model=model_name, temperature=0.3, org_id=org_id)
raw = await llm_call_msgs(msgs, model=model_name, temperature=0.3, org_id=org_id, project_id=project_id, session_id='task:%s' % task_id)
except Exception as e:
err_msg = f"{type(e).__name__}: {str(e)[:400]}"
from .task_capability import mark_failed
@ -3407,7 +3408,7 @@ async def qc_review_run(project_id, agent_id=None, model_name=None):
msgs.append({"role": "user", "content":
"已检查足够信息。现在必须立即输出 review_approve 或 review_reject禁止再调用其它工具。"})
try:
raw = await llm_call_msgs(msgs, model=model_name, temperature=0.2, org_id=org_id)
raw = await llm_call_msgs(msgs, model=model_name, temperature=0.2, org_id=org_id, project_id=project_id, session_id='task:%s' % task_id)
except Exception as e:
err_msg = f"{type(e).__name__}: {str(e)[:400]}"
from .task_capability import mark_failed
@ -3606,7 +3607,7 @@ async def retrospective_run(project_id, agent_id=None, model_name=None):
"已执行足够轮次。现在必须立即收尾write_file 写复盘报告(若未写),然后输出 deliver 提交,禁止再调其它工具。"})
try:
raw = await asyncio.wait_for(
llm_call_msgs(msgs, model=model_name, temperature=0.3, org_id=org_id),
llm_call_msgs(msgs, model=model_name, temperature=0.3, org_id=org_id, project_id=project_id, session_id='task:%s' % task_id),
timeout=_LLM_HARD_TIMEOUT)
except Exception as e:
err_msg = f"{type(e).__name__}: {str(e)[:400]}"

View File

@ -717,6 +717,7 @@ class AgentExecutor:
temperature=0,
org_id=self.org_id,
purpose='utility',
session_id=self.session_id,
timeout=self._UTILITY_TIMEOUT,
)
m = _re.search(r"\[[^\]]*\]", content or "")
@ -832,6 +833,7 @@ class AgentExecutor:
model=self.model_name,
temperature=self.config.temperature,
org_id=self.org_id,
session_id=self.session_id,
)
except Exception as e:
# 原生视觉降级2026-09-10模型不支持多模态 content 数组时
@ -842,7 +844,8 @@ class AgentExecutor:
try:
return await llm_call_msgs_native(
self._msgs, tools=tools_schema, model=self.model_name,
temperature=self.config.temperature, org_id=self.org_id)
temperature=self.config.temperature, org_id=self.org_id,
session_id=self.session_id)
except Exception:
pass
logger.error(f"native function calling failed, fallback to text: {e}")
@ -855,12 +858,14 @@ class AgentExecutor:
model=self.model_name,
temperature=self.config.temperature,
org_id=self.org_id,
session_id=self.session_id,
)
except Exception as e:
if self._degrade_images_if_needed(e):
content = await llm_call_msgs(
self._msgs, model=self.model_name,
temperature=self.config.temperature, org_id=self.org_id)
temperature=self.config.temperature, org_id=self.org_id,
session_id=self.session_id)
else:
raise
return {"content": content or "", "tool_calls": []}
@ -1238,7 +1243,7 @@ class AgentExecutor:
pnames = [getattr(r, "name", "") for r in (all_recs or [])]
classify_prompt = f"用户输入: {name}\n项目列表: {', '.join(pnames)}\n\n判断用户想要哪个项目。只回复项目名或\"不存在\""
matched = await llm_call(classify_prompt, temperature=0.0, org_id=self.org_id,
purpose='utility')
purpose='utility', session_id=self.session_id)
matched = matched.strip().strip('"').strip("'")
recs2 = await sor.sqlExe(
@ -1628,7 +1633,7 @@ class AgentExecutor:
from .platform_model_tools import tool_invoke_model
return await tool_invoke_model(
p, self.org_id or "0", user_id=self.user_id or "",
project_id=self.project_id or "")
project_id=self.project_id or "", session_id=self.session_id or "")
async def _t_web_search(self, sor, p, pid):
"""联网检索(薄壳委托 web_tools 唯一实现SSRF 防护在其中)。"""
@ -2190,6 +2195,7 @@ class AgentExecutor:
model=self.model_name,
org_id=self.org_id,
purpose='utility',
session_id=self.session_id,
timeout=self._UTILITY_TIMEOUT,
)
return summary[:500]

View File

@ -41,9 +41,14 @@ async def _self_base_url():
return "http://127.0.0.1:%d/pipeline-llm/api/v1" % port
async def _get_internal_token(org_id, user_id, model_name):
"""取/发内部短期 token。失败抛 ValueError消息真实可行动"""
key = (org_id or '0', user_id or '', model_name or '')
async def _get_internal_token(org_id, user_id, model_name, project_id=''):
"""取/发内部短期 token。失败抛 ValueError消息真实可行动
project_id2026-09-10 用户定夺 token pipeline_llm_tokens.project_id
推理端点透传到 llm_usage.project_id 支撑按项目统计费用缓存键必须含
project_id否则同 org/user/model 的不同项目会串用同一 token费用记错项目
"""
key = (org_id or '0', user_id or '', model_name or '', project_id or '')
now = time.time()
ent = _token_cache.get(key)
if ent and ent['calls'] < _TOKEN_MAX_LOCAL_CALLS \
@ -52,7 +57,7 @@ async def _get_internal_token(org_id, user_id, model_name):
return ent['token']
from .llm_proxy import create_llm_token
ok, token = await create_llm_token(
org_id or '0', project_id='', task_id='', model_name=model_name or '',
org_id or '0', project_id=project_id or '', task_id='', model_name=model_name or '',
purpose='internal_bridge', ttl_hours=_TOKEN_TTL_HOURS,
max_calls=500, created_by=user_id or '')
if not ok:
@ -64,7 +69,8 @@ async def _get_internal_token(org_id, user_id, model_name):
return token
async def _http_chat(payload, org_id, user_id, model_name, timeout: int = 0):
async def _http_chat(payload, org_id, user_id, model_name, timeout: int = 0,
project_id: str = ''):
"""POST 本进程推理端点。返回上游响应 dict失败抛 ValueError消息真实可行动
timeout客户端等待秒数0=缺省 330异步生成模型视频等端点侧最长
@ -72,7 +78,7 @@ async def _http_chat(payload, org_id, user_id, model_name, timeout: int = 0):
"""
import aiohttp
token = await _get_internal_token(org_id, user_id, model_name)
token = await _get_internal_token(org_id, user_id, model_name, project_id)
base = await _self_base_url()
# ⚠️ "Bearer " 前缀用拼接构造——字面量写在源码里会被脱敏工具替换成 ***
# 2026-09-04 实测Authorization 头变成 "***plk-..." 致端点校验失败)
@ -111,13 +117,16 @@ def _no_llm_error(model_name=None, org_id=None) -> ValueError:
async def llm_call(prompt: str, model: str = None, temperature: float = 0.7,
org_id: str = None, user_id: str = None, purpose: str = '',
timeout: int = 0) -> str:
timeout: int = 0, project_id: str = '', session_id: str = '') -> str:
"""Call LLM and return text response.
统一走模型治理推理 API门禁链 + 双维度记账
org_id 为空 = 系统级'0'与旧语义不过滤机构等价
purpose='utility'辅助任务分类/选择/摘要治理层按用途选模型链
机构策略配置辅助模型优先 payload _purpose 键透传到端点
session_id2026-09-11 用户定夺会话粘性键同会话固定同一账号
上游 KV/前缀缓存按账号隔离固定账号命中缓存省钱 payload._session_id
透传到端点不进 token 缓存键token 与账号选择解耦
"""
# 兼容旧优先级harnessed_agent若宿主加载了独立推理后端
try:
@ -141,7 +150,10 @@ async def llm_call(prompt: str, model: str = None, temperature: float = 0.7,
payload["_purpose"] = purpose
if timeout:
payload["_timeout"] = int(timeout)
data = await _http_chat(payload, org_id or '0', user_id or '', model or '')
if session_id:
payload["_session_id"] = session_id
data = await _http_chat(payload, org_id or '0', user_id or '', model or '',
project_id=project_id or '')
try:
return data["choices"][0]["message"]["content"]
except (KeyError, IndexError, TypeError) as e:
@ -156,18 +168,22 @@ async def call_llm(tenant_id: str, prompt: str, model: str = None, temperature:
async def llm_call_msgs(messages: list, model: str = None, temperature: float = 0.7,
org_id: str = None, user_id: str = None, purpose: str = '',
timeout: int = 0) -> str:
timeout: int = 0, project_id: str = '', session_id: str = '') -> str:
"""Call LLM with full message array (system/user/assistant).
purpose='utility'辅助任务分类/选择/摘要治理层按用途选模型链
timeout单次上游调用超时秒数0=用端点默认上限 900超长文本提取用
session_id2026-09-11会话粘性账号键 payload._session_id 透传
"""
payload = {"model": model or '', "messages": messages, "temperature": temperature}
if purpose:
payload["_purpose"] = purpose
if timeout:
payload["_timeout"] = int(timeout)
data = await _http_chat(payload, org_id or '0', user_id or '', model or '')
if session_id:
payload["_session_id"] = session_id
data = await _http_chat(payload, org_id or '0', user_id or '', model or '',
project_id=project_id or '')
try:
return data["choices"][0]["message"]["content"]
except (KeyError, IndexError, TypeError) as e:
@ -177,13 +193,15 @@ async def llm_call_msgs(messages: list, model: str = None, temperature: float =
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, purpose: str = '') -> dict:
user_id: str = None, purpose: str = '',
project_id: str = '', session_id: str = '') -> dict:
"""Native function calling. 传入 tools JSON schema返回 message dict。
Returns:
{"content": str, "tool_calls": [{"id","type","function":{"name","arguments"}}]}
当模型返回 tool_calls content 通常为空字符串
purpose='utility'辅助任务分类/选择/摘要治理层按用途选模型链
session_id2026-09-11会话粘性账号键 payload._session_id 透传
"""
payload = {"model": model or '', "messages": messages, "temperature": temperature}
if tools:
@ -191,7 +209,10 @@ async def llm_call_msgs_native(messages: list, tools: list = None, model: str =
payload["tool_choice"] = "auto"
if purpose:
payload["_purpose"] = purpose
data = await _http_chat(payload, org_id or '0', user_id or '', model or '')
if session_id:
payload["_session_id"] = session_id
data = await _http_chat(payload, org_id or '0', user_id or '', model or '',
project_id=project_id or '')
try:
msg = data["choices"][0]["message"]
except (KeyError, IndexError, TypeError) as e:
@ -204,7 +225,8 @@ async def llm_call_msgs_native(messages: list, tools: list = None, model: str =
async def llm_infer(payload: dict, model: str = None, org_id: str = None,
user_id: str = None, timeout: int = 0) -> dict:
user_id: str = None, timeout: int = 0,
project_id: str = '', session_id: str = '') -> dict:
"""通用推理全能力2026-09-07透传任意 payload 到统一推理端点,
返回上游响应 dictOpenAI 兼容 choices生成类另带 media/output/task_id
@ -222,7 +244,9 @@ async def llm_infer(payload: dict, model: str = None, org_id: str = None,
body["model"] = model
if timeout:
body["_timeout"] = int(timeout)
if session_id and not body.get('_session_id'):
body["_session_id"] = session_id
# 客户端等待须覆盖端点侧预算_timeout 端点封顶 900+ 余量,否则客户端先断连
client_timeout = min(int(timeout or 0), 900) + 60 if timeout else 0
return await _http_chat(body, org_id or '0', user_id or '', model or '',
timeout=client_timeout)
timeout=client_timeout, project_id=project_id or '')

View File

@ -214,7 +214,8 @@ def load_model_skill(model_name):
async def judge_completeness(model_name, capability, task, biz,
description="", org_id="0", user_id=""):
description="", org_id="0", user_id="",
project_id=""):
"""完备性判断主体LLM 裁决2026-09-08 用户定夺)。
模型有配套技能 按技能输入契约判断 无技能 按能力语义判断
@ -266,7 +267,7 @@ async def judge_completeness(model_name, capability, task, biz,
raw = await llm_call_msgs(
[{"role": "user", "content": prompt}],
temperature=0, org_id=org_id or "0", user_id=user_id or "",
purpose="utility", timeout=60)
purpose="utility", timeout=60, project_id=project_id or "")
except Exception as e:
logger.warning("judge_completeness LLM 失败(走代码表兜底): %s", repr(e)[:150])
return _check_completeness_code(cap, task, biz)
@ -338,7 +339,7 @@ async def _resolve_capability(model_name, org_id):
return "", ""
async def tool_invoke_model(params, org_id, user_id=""):
async def tool_invoke_model(params, org_id, user_id="", project_id="", session_id=""):
"""调用平台模型完成生成类任务(文生图/视频/语音等非对话能力)。
流程解析 task/model/capability/params 未指定 model 时按 task 自动
@ -404,7 +405,8 @@ async def tool_invoke_model(params, org_id, user_id=""):
else:
gap = await judge_completeness(model, cap_gate, task, biz,
description=desc, org_id=org_id or "0",
user_id=user_id or "")
user_id=user_id or "",
project_id=project_id or "")
cap = cap_gate
if gap:
logger.info("invoke_model 完备性门禁拦截(%s: %s", cap or "?", gap[:80])
@ -417,7 +419,8 @@ async def tool_invoke_model(params, org_id, user_id=""):
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)
user_id=user_id or "", timeout=600, project_id=project_id or "",
session_id=session_id or "")
except Exception as e:
return "FAIL: 模型「" + model + "」调用失败:" + str(e)[:400]
@ -475,10 +478,12 @@ PLATFORM_MODEL_TOOLS_V1 = [
]
async def exec_platform_model_tool(tool, params, org_id, user_id=""):
async def exec_platform_model_tool(tool, params, org_id, user_id="",
project_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 await tool_invoke_model(params, org_id, user_id=user_id,
project_id=project_id)
return "未实现: " + str(tool)