fix: 修复 5 项确定 bug + Provider UX 重做 + 文档统一
Bug fixes: - fix(dao): AsyncSession.delete 补齐漏掉的 await(provider/user/individual 共 4 处) - fix(worker): result.data.output → result.output.output(pydantic-ai 1.x API 适配) - fix(api): 删除 create_worker_from_template 死端点(ORM 字段不匹配必崩) - fix(api): /provider/test 按 provider_type 分支适配 Anthropic/Gemini/OpenAI 三种协议 - fix(chat): SSE 流式聊天在 distributed 模式 fallback 到非流式,避免 asyncio.Queue 序列化崩溃 Features (previously unstaged): - feat(provider): Provider 管理页重做(品牌图标、5 种类型、Test Connection、编辑模式) - feat(provider): 新增 Gemini provider_type 支持 - feat(workflow): Finalize 节点输出 blackboard 摘要 + 失败原因;步骤完成/失败实时推送 SSE - feat(i18n): regulatory_node 提示词从路由模式改为直接对话模式(中英双语) - feat(consciousness): dynamic_prompt 支持 locale 国际化 - feat(logs): SystemLogsView 自动刷新 + 暂停按钮 Docs: - docs: README/README-EN 统一为"开源通用多 Agent 协作平台"口径 - docs: ROADMAP 按 v0.1.x / v0.2.x / v0.3.x 重组 - docs: project.md 重写为结构化项目介绍 Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
+17
-6
@@ -163,15 +163,15 @@ async def stream_chat_message(
|
||||
request: Request,
|
||||
token_data: TokenData = Depends(Accessor.get_current_user),
|
||||
):
|
||||
"""SSE 流式聊天端点:通过 regulatory_node agent 流式输出,支持工具调用。"""
|
||||
"""SSE 流式聊天端点:standalone 模式下逐 token 流式输出;distributed 模式 fallback 到整段回复。"""
|
||||
from kilostar.utils.standalone_proxy import _STANDALONE
|
||||
|
||||
postgres_database = ray_actor_hook("postgres_database").postgres_database
|
||||
|
||||
# 存用户消息
|
||||
await postgres_database.add_chat_message.remote(
|
||||
chat_id=chat_id, message=request_body.message, message_owner="user"
|
||||
)
|
||||
|
||||
# 构造 MessageRequest payload
|
||||
payload = MessageRequest(
|
||||
platform="client",
|
||||
user_name=token_data.user_id,
|
||||
@@ -180,9 +180,21 @@ async def stream_chat_message(
|
||||
)
|
||||
|
||||
regulatory_node = ray_actor_hook("regulatory_node").regulatory_node
|
||||
token_queue = asyncio.Queue()
|
||||
|
||||
# stream_working.remote() returns an asyncio.Task in standalone mode
|
||||
if not _STANDALONE:
|
||||
async def fallback_generator():
|
||||
resp = await regulatory_node.working.remote(payload)
|
||||
full_response = resp.reply_message if resp else ""
|
||||
if full_response:
|
||||
await postgres_database.add_chat_message.remote(
|
||||
chat_id=chat_id, message=full_response, message_owner="regulatory_node"
|
||||
)
|
||||
yield f"data: {json.dumps({'token': full_response})}\n\n"
|
||||
yield f"data: {json.dumps({'done': True, 'full_message': full_response})}\n\n"
|
||||
|
||||
return StreamingResponse(fallback_generator(), media_type="text/event-stream")
|
||||
|
||||
token_queue = asyncio.Queue()
|
||||
stream_task = regulatory_node.stream_working.remote(payload, token_queue)
|
||||
|
||||
async def event_generator():
|
||||
@@ -207,7 +219,6 @@ async def stream_chat_message(
|
||||
full_response = "抱歉,生成回复时出错。"
|
||||
yield f"data: {json.dumps({'token': full_response})}\n\n"
|
||||
|
||||
# 流结束,存入数据库
|
||||
if full_response:
|
||||
await postgres_database.add_chat_message.remote(
|
||||
chat_id=chat_id,
|
||||
|
||||
@@ -27,7 +27,7 @@ provider_router = APIRouter(prefix="/api/v1/provider", tags=["provider"])
|
||||
class ProviderRegister(BaseModel):
|
||||
"""``POST /provider`` 入参:注册一个模型 Provider 的最小字段集。"""
|
||||
|
||||
provider_type: Literal["openai", "claude", "deepseek"]
|
||||
provider_type: Literal["openai", "claude", "deepseek", "gemini"]
|
||||
provider_title: str
|
||||
provider_url: str
|
||||
provider_apikey: str
|
||||
@@ -72,6 +72,63 @@ async def get_provider_list(
|
||||
return {"provider_list": masked}
|
||||
|
||||
|
||||
@provider_router.post("/test")
|
||||
async def test_provider_connection(
|
||||
provider_register: ProviderRegister,
|
||||
_: TokenData = Depends(Accessor.get_current_user),
|
||||
) -> Dict[str, Any]:
|
||||
"""测试 Provider 连接:按 provider_type 选择对应协议拉取模型列表。"""
|
||||
import httpx
|
||||
|
||||
ptype = provider_register.provider_type
|
||||
url = provider_register.provider_url
|
||||
apikey = provider_register.provider_apikey
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=10.0) as client:
|
||||
if ptype == "claude":
|
||||
endpoint = f"{url}/v1/models"
|
||||
headers = {
|
||||
"x-api-key": apikey,
|
||||
"anthropic-version": "2023-06-01",
|
||||
}
|
||||
response = await client.get(endpoint, headers=headers)
|
||||
if response.status_code == 200:
|
||||
data = response.json()
|
||||
models = [m["id"] for m in data.get("data", [])]
|
||||
return {"success": True, "models": sorted(models), "model_count": len(models)}
|
||||
return {"success": False, "error": f"HTTP {response.status_code}", "models": []}
|
||||
|
||||
elif ptype == "gemini":
|
||||
endpoint = f"{url}/models"
|
||||
params = {"key": apikey}
|
||||
response = await client.get(endpoint, params=params)
|
||||
if response.status_code == 200:
|
||||
data = response.json()
|
||||
models = [m.get("name", "").removeprefix("models/") for m in data.get("models", [])]
|
||||
return {"success": True, "models": sorted(models), "model_count": len(models)}
|
||||
return {"success": False, "error": f"HTTP {response.status_code}", "models": []}
|
||||
|
||||
else:
|
||||
if "/v1" not in url:
|
||||
endpoint = f"{url}/v1/models"
|
||||
else:
|
||||
endpoint = f"{url}/models"
|
||||
headers = {
|
||||
"Authorization": f"Bearer {apikey}",
|
||||
"Content-Type": "application/json",
|
||||
}
|
||||
response = await client.get(endpoint, headers=headers)
|
||||
if response.status_code == 200:
|
||||
data = response.json()
|
||||
models = [m["id"] for m in data.get("data", [])]
|
||||
return {"success": True, "models": sorted(models), "model_count": len(models)}
|
||||
return {"success": False, "error": f"HTTP {response.status_code}", "models": []}
|
||||
|
||||
except Exception as e:
|
||||
return {"success": False, "error": str(e), "models": []}
|
||||
|
||||
|
||||
@provider_router.delete("/{provider_title}")
|
||||
async def delete_provider(
|
||||
provider_title: str,
|
||||
|
||||
Reference in New Issue
Block a user