deerflow-code/newpython/OUT.py
2026-09-07 18:24:55 +08:00

483 lines
17 KiB
Python
Raw Permalink 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.

"""
测试流式接口的完整脚本
API流程:
1. 登录获取token
2. 创建thread(可选,也可以使用无状态接口)
3. 调用流式接口
两种方式:
- 方式1: POST /api/runs/stream (无状态,自动创建thread)
- 方式2: POST /api/threads/{thread_id}/runs/stream (需要先创建thread)
"""
import requests
import json
import asyncio
import httpx # 异步 HTTP 客户端
# ============ 配置区域 ============
BASE_URL = "http://127.0.0.1:8001"
LOGIN_URL = f"{BASE_URL}/api/v1/auth/login/username"
THREADS_URL = f"{BASE_URL}/api/threads"
RUNS_STREAM_URL = f"{BASE_URL}/api/runs/stream" # 无状态流式接口
# ============ 登录获取token ============
def get_token(username: str) -> str:
"""登录获取token"""
response = requests.post(
LOGIN_URL,
json={"username": username},
headers={"Content-Type": "application/json"}
)
# print(f"登录状态码: {response.status_code}")
data = response.json()
# print(f"登录响应: {json.dumps(data, indent=2, ensure_ascii=False)}")
# 根据实际返回结构调整
token = data.get("token") or data.get("access_token") or data.get("data", {}).get("token")
return token
# ============ 方式2: 先创建thread,再调用流式接口 ============
def create_thread(token: str) -> str:
"""创建thread,返回thread_id"""
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {token}"
}
response = requests.post(
THREADS_URL,
json={}, # 空body,使用默认值
headers=headers
)
print(f"创建thread状态码: {response.status_code}")
data = response.json()
print(f"创建thread响应: {json.dumps(data, indent=2, ensure_ascii=False)}")
return data.get("thread_id")
def get_agent_name_description(agent_name="Taiwan-intelligence-gathering-agent"):
url = f"{BASE_URL}/api/agents/{agent_name}"
token=get_token("ji")
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {token}"
}
response = requests.get(
url,
headers=headers,
)
print(response.json())
return response.json()["name"].lower(),"agent_name:"+response.json()["name"].lower()+","+response.json()["description"]
# ============ 创建智能体 ============
def create_agent(
name: str,
description: str = "",
model: str | None = None,
tool_groups: list[str] | None = None,
skills: list[str] | None = None,
soul: str = "",
):
"""
创建一个新的智能体
Args:
name: 智能体名称(必须是 hyphen-case,如 my-agent)
description: 智能体描述
model: 可选的模型覆盖(如 gpt-4, claude-3-opus 等)
tool_groups: 工具组白名单(如 ["web", "file"])
skills: 技能白名单(如 ["skill1", "skill2"])
soul: SOUL.md 内容(智能体人格/身份描述)
Returns:
创建结果
"""
TOKEN =get_token("ji")
url = f"{BASE_URL}/api/agents"
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {TOKEN}"
}
payload = {
"name": name,
"description": description,
"soul": soul
}
# 添加可选字段
if model is not None:
payload["model"] = model
if tool_groups is not None:
payload["tool_groups"] = tool_groups
if skills is not None:
payload["skills"] = skills
print(f"请求: POST {url}")
print(f"请求体: {json.dumps(payload, indent=2, ensure_ascii=False)}")
response = requests.post(url, json=payload, headers=headers)
print(f"状态码: {response.status_code}")
if response.status_code == 201:
result = response.json()
print(f"✓ 智能体创建成功!")
print(f"响应: {json.dumps(result, indent=2, ensure_ascii=False)}")
return result
else:
print(f"✗ 创建失败: {response.text}")
return f"✗ 创建失败: {response.text}"
def create_main_agent(zhihui_name,agent_names):
d={}
token=get_token("ji")
l=[]
for i in agent_names:
# print(i)
name,description=get_agent_name_description(i)
thread_id_tmp=create_thread(token)
d[name]=thread_id_tmp
l.append(description)
name="你是总控控制协调智能体,负责协调其他智能体干活"
miaoshu=f"""#核心定位:
用途:跟用用户的任务目标,调用其他智能体干活,只负责协调汇总,当你觉得已有信息支持回答用户问题时候,停止调度,并输出最后的结果
#调度方法:
通过使用agent_orchestration技能调用其他智能体干活,使用技能时传入agent_name和需要做的活
#可调度智能体:
{l}"""
thread_id_tmp=create_thread(token)
engname=zhihui_name
d[engname]=thread_id_tmp
tmp_skills=["agent_orchestration"]
tmpres=create_agent(
name=engname,
description=name,
model="glm-5", # 指定模型
skills=tmp_skills, # 关联技能
soul=miaoshu)
# print("创建智能体响应:",tmpres)
if "失败" in tmpres and "exists" not in tmpres:
return tmpres
print(f"创建协调智能体成功,可以协调列表:{l}")
return d
import re
def extract_name_task(text="python3 scripts/agent_orchestration.py \"Us-intelligence-gathering-agent\" \"收集和分析美国人民对特朗普的看法,包括支持率和反对率、不同群体的观点差异、主要支持理由和反对理由等\""):
try:
pattern = r'\.py\s+"([^"]*)"\s+"([^"]*)"'
match = re.search(pattern, text)
if match:
agent_name = match.group(1)
task_content = match.group(2)
return agent_name,task_content
except:
return "",""
def update_agent_message(thread_id,new_content):
token=get_token("ji")
url = f"{BASE_URL}/api/threads/{thread_id}/state"
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {token}"
}
tmp = requests.get(
url,
headers=headers
).json()
# print(tmp)
if "messages" in tmp["values"]:
tmp["values"]["messages"].append({"type": "human", "content": new_content})
else:
tmp["values"]["messages"]=[]
tmp["values"]["messages"].append({"type": "human", "content": new_content})
url = f"{BASE_URL}/api/threads/{thread_id}/state"
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {token}"
}
response = requests.post(
url,
json=tmp,
headers=headers,
stream=True
)
del tmp
print(f"线程{thread_id}更新状态码: {response.status_code}")
def Broadcast_message(message,d_agent_thread_id,pingbi_agent_name):
for key in d_agent_thread_id:
if key not in pingbi_agent_name:
update_agent_message(d_agent_thread_id[key],message)
print("广播消息成功:",key,message)
null=None
true=True
false=False
async def leader_agent(d_agent_thread_id, agent_name, new_message, skill_stop_names):
"""使用 httpx 异步请求的 leader_agent"""
# 使用 to_thread 异步运行同步的 Broadcast_message
await asyncio.to_thread(Broadcast_message, "总控智能体收到用户消息:" + new_message, d_agent_thread_id, [agent_name])
thread_id = d_agent_thread_id[agent_name]
l = []
token = get_token("ji")
url = f"{BASE_URL}/api/threads/{thread_id}/runs/stream"
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {token}"
}
payload = {
"input": {
"messages": [
{
"type": "human",
"content": [
{
"type": "text",
"text": new_message
}
],
"additional_kwargs": {}
}
]
},
"config": {
"recursion_limit": 1000
},
"context": {
"agent_name": agent_name,
"model_name": "deepseek-chat",
"mode": "pro",
"reasoning_effort": "medium",
"thinking_enabled": true,
"is_plan_mode": false,
"subagent_enabled": false,
"thread_id": thread_id
},
"stream_mode": [
"messages-tuple",
"values",
],
"stream_subgraphs": true,
"stream_resumable": true,
"assistant_id": "lead_agent",
"on_disconnect": "continue",
"excluded_tools": ["web_search"],
"skill_stop_names": skill_stop_names,
}
async with httpx.AsyncClient(timeout=None) as client:
async with client.stream("POST", url, json=payload, headers=headers) as response:
print(f"状态码: {response.status_code}")
if response.status_code != 200:
body = await response.aread()
raise RuntimeError(
f"upstream {url} returned HTTP {response.status_code}: "
f"{body.decode('utf-8', 'replace')[:1000]}"
)
async for line in response.aiter_lines():
if line:
yield line
if line[:5] == "data:":
print(line)
l.append(line)
if len(l) < 2:
raise RuntimeError(
f"upstream stream produced no usable data lines (got {len(l)} line(s)). "
f"Check the server log on {url} for the real failure."
)
null = None
tmpres = eval(l[-2][5:])
tmpres = tmpres["messages"][-1]
tmp_task = []
if tmpres.get("tool_calls", []) == []:
yield {"status": []}
else:
for tool_call in tmpres.get("tool_calls", []):
try:
agentname, task_text = extract_name_task(str(tool_call["args"]))
print("agentname", agentname)
print("task_text:", task_text)
# 使用 to_thread 异步运行同步的 Broadcast_message
await asyncio.to_thread(Broadcast_message, f"总控智能体给{agentname}派活,内容为:" + task_text, d_agent_thread_id, [agent_name, agentname])
tmp_task.append([agentname, task_text])
except:
pass
yield {"status": tmp_task}
async def special_agent(d_agent_thread_id, agent_name, new_message, skill_stop_names):
"""使用 httpx 异步请求的 special_agent"""
thread_id = d_agent_thread_id[agent_name]
l = []
token = get_token("ji")
url = f"{BASE_URL}/api/threads/{thread_id}/runs/stream"
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {token}"
}
payload = {
"input": {
"messages": [
{
"type": "human",
"content": [
{
"type": "text",
"text": new_message
}
],
"additional_kwargs": {}
}
]
},
"config": {
"recursion_limit": 1000
},
"context": {
"agent_name": agent_name,
"model_name": "deepseek-chat",
"mode": "pro",
"reasoning_effort": "medium",
"thinking_enabled": true,
"is_plan_mode": false,
"subagent_enabled": false,
"thread_id": thread_id
},
"stream_mode": [
"messages-tuple",
"values",
],
"stream_subgraphs": true,
"stream_resumable": true,
"assistant_id": "lead_agent",
"on_disconnect": "continue",
"excluded_tools": ["web_search", "ask_clarification", "present_files", "view_image"],
"skill_stop_names": skill_stop_names,
}
async with httpx.AsyncClient(timeout=None) as client:
async with client.stream("POST", url, json=payload, headers=headers) as response:
print(f"状态码: {response.status_code}")
if response.status_code != 200:
body = await response.aread()
raise RuntimeError(
f"upstream {url} returned HTTP {response.status_code}: "
f"{body.decode('utf-8', 'replace')[:1000]}"
)
async for line in response.aiter_lines():
if line:
yield line
if line[:5] == "data:":
print(line)
l.append(line)
if len(l) < 2:
raise RuntimeError(
f"upstream stream produced no usable data lines (got {len(l)} line(s)). "
f"Check the server log on {url} for the real failure."
)
null = None
tmpres = eval(l[-2][5:])
tmpres = tmpres["messages"][-1]
# 使用 to_thread 异步运行同步的 Broadcast_message
await asyncio.to_thread(Broadcast_message, f"子智能体{agent_name}完成{new_message}工作,交付内容为:" + tmpres["content"], d_agent_thread_id, [agent_name])
yield {"status": f"子智能体{agent_name}完成{new_message}工作"}
async def agent_init(agent_names, main_agent_name):
"""异步的 agent_init,同步函数调用使用 to_thread"""
agent_names = [i.lower() for i in agent_names]
# 使用 to_thread 异步运行同步的 create_main_agent
tmpd = await asyncio.to_thread(create_main_agent, main_agent_name, agent_names)
message = """
你的职责:你只是若干个智能体的其中一员,你们的共同目标是完成一个庞大的任务,总控控制协调智能体会分配特定的活给你,你只需要用你自身的能力来完成这个特定的活
注意:你只是一个子智能体,你不必把你的分析结果存放到文件里面去,你只关心于完成总控控制协调智能体分配的活即可
"""
# 使用 to_thread 异步运行同步的 Broadcast_message
await asyncio.to_thread(Broadcast_message, message, tmpd, [main_agent_name])
return main_agent_name, tmpd
from fastapi import FastAPI, HTTPException,Body
from fastapi.responses import StreamingResponse
from fastapi.middleware.cors import CORSMiddleware
import asyncio
import uvicorn
app = FastAPI()
# CORS 配置
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
import time
@app.post("/generate_mul_agent")
async def generate_stream(payload: dict = Body(...)):
"""
根据 payload 中的 agent_type 字段判断调用 leader_agent 还是 special_agent
payload 结构示例:
{
"agent_type": "leader", # 或 "special"
"d_agent_thread_id": {...},
"agent_name": "xxx",
"new_message": "xxx",
"skill_stop_names": [...]
}
"""
# 从 payload 中提取参数
agent_type = payload.get("agent_type", "leader") # 默认为 leader
d_agent_thread_id = payload.get("d_agent_thread_id", {})
agent_name = payload.get("agent_name", "")
new_message = payload.get("new_message", "")
skill_stop_names = payload.get("skill_stop_names", [])
# 流式响应生成器
async def stream_generator():
# 根据 agent_type 选择调用的函数
if agent_type == "leader":
async_generator = leader_agent(d_agent_thread_id, agent_name, new_message, skill_stop_names)
else: # special
async_generator = special_agent(d_agent_thread_id, agent_name, new_message, skill_stop_names)
# 异步迭代生成器
async for chunk in async_generator:
if isinstance(chunk, dict):
yield f"data: {json.dumps(chunk, ensure_ascii=False)}\n\n"
else:
yield f"{chunk}"
# 使用 SSE 流式响应
return StreamingResponse(
stream_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no",
"Connection": "keep-alive",
}
)
@app.post("/agent_init")
async def init_agents(payload: dict = Body(...)):
"""
初始化智能体接口
payload 结构示例:
{
"agent_names": ["agent1", "agent2"],
"main_agent_name": "main-agent"
}
"""
agent_names = payload.get("agent_names", [])
main_agent_name = payload.get("main_agent_name", "main-agent")
# 调用异步的 agent_init
name, thread_id_dict = await agent_init(agent_names, main_agent_name)
return {
"status": "success",
"main_agent_name": name,
"thread_ids": thread_id_dict
}
if __name__ == "__main__":
uvicorn.run("testapi:app", host="0.0.0.0", port=8095, workers=1,limit_concurrency=1000) # 多进程提升并发能力