diff --git a/pipeline_service/agent_loop.py b/pipeline_service/agent_loop.py index e66cd16..64ce1d3 100644 --- a/pipeline_service/agent_loop.py +++ b/pipeline_service/agent_loop.py @@ -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]}" diff --git a/pipeline_service/agent_loop_v2.py b/pipeline_service/agent_loop_v2.py index 7d1c5f6..6bc4e3f 100644 --- a/pipeline_service/agent_loop_v2.py +++ b/pipeline_service/agent_loop_v2.py @@ -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] diff --git a/pipeline_service/llm_bridge.py b/pipeline_service/llm_bridge.py index 6fb5f5d..548abed 100644 --- a/pipeline_service/llm_bridge.py +++ b/pipeline_service/llm_bridge.py @@ -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_id(2026-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_id(2026-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_id(2026-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_id(2026-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 到统一推理端点, 返回上游响应 dict(OpenAI 兼容 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 '') diff --git a/pipeline_service/platform_model_tools.py b/pipeline_service/platform_model_tools.py index 992361d..164c145 100644 --- a/pipeline_service/platform_model_tools.py +++ b/pipeline_service/platform_model_tools.py @@ -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)