schedule 通过 aiohttp POST 请求调用
http://127.0.0.1:8001/api/v2/chat/completions
This commit is contained in:
parent
4090b4d734
commit
393c4e4138
@ -9,6 +9,8 @@ import asyncio
|
||||
import logging
|
||||
import time
|
||||
import yaml
|
||||
import aiohttp
|
||||
import json
|
||||
from datetime import datetime, timezone, timedelta
|
||||
from pathlib import Path
|
||||
from typing import Optional
|
||||
@ -136,7 +138,7 @@ class ScheduleExecutor:
|
||||
logger.info(f"Executing scheduled task: {task_id} ({task.get('name', '')}) for bot={bot_id} user={user_id}")
|
||||
|
||||
# 调用 agent
|
||||
response_text = await self._call_agent(bot_id, user_id, task)
|
||||
response_text = await self._call_agent_v2(bot_id, user_id, task)
|
||||
|
||||
# 写入日志
|
||||
duration_ms = int((time.time() - start_time) * 1000)
|
||||
@ -156,42 +158,35 @@ class ScheduleExecutor:
|
||||
finally:
|
||||
self._executing_tasks.discard(task_id)
|
||||
|
||||
async def _call_agent(self, bot_id: str, user_id: str, task: dict) -> str:
|
||||
"""构建 AgentConfig 并调用 agent"""
|
||||
from routes.chat import create_agent_and_generate_response
|
||||
from utils.fastapi_utils import fetch_bot_config, create_project_directory
|
||||
from agent.agent_config import AgentConfig
|
||||
async def _call_agent_v2(self, bot_id: str, user_id: str, task: dict) -> str:
|
||||
"""通过 HTTP 调用 /api/v2/chat/completions 接口"""
|
||||
from utils.fastapi_utils import generate_v2_auth_token
|
||||
|
||||
bot_config = await fetch_bot_config(bot_id)
|
||||
project_dir = create_project_directory(
|
||||
bot_config.get("dataset_ids", []),
|
||||
bot_id,
|
||||
bot_config.get("skills")
|
||||
)
|
||||
url = f"http://127.0.0.1:8001/api/v2/chat/completions"
|
||||
auth_token = generate_v2_auth_token(bot_id)
|
||||
|
||||
messages = [{"role": "user", "content": task["message"]}]
|
||||
payload = {
|
||||
"messages": [{"role": "user", "content": task["message"]}],
|
||||
"stream": False,
|
||||
"bot_id": bot_id,
|
||||
"tool_response": False,
|
||||
"session_id": f"schedule_{task['id']}",
|
||||
"user_identifier": user_id,
|
||||
}
|
||||
|
||||
config = AgentConfig(
|
||||
bot_id=bot_id,
|
||||
api_key=bot_config.get("api_key"),
|
||||
model_name=bot_config.get("model", "qwen3-next"),
|
||||
model_server=bot_config.get("model_server", ""),
|
||||
language=bot_config.get("language", "ja"),
|
||||
system_prompt=bot_config.get("system_prompt"),
|
||||
mcp_settings=bot_config.get("mcp_settings", []),
|
||||
user_identifier=user_id,
|
||||
session_id=f"schedule_{task['id']}",
|
||||
project_dir=project_dir,
|
||||
stream=False,
|
||||
tool_response=False,
|
||||
messages=messages,
|
||||
dataset_ids=bot_config.get("dataset_ids", []),
|
||||
enable_memori=bot_config.get("enable_memory", False),
|
||||
shell_env=bot_config.get("shell_env") or {},
|
||||
)
|
||||
headers = {
|
||||
"Content-Type": "application/json",
|
||||
"Authorization": f"Bearer {auth_token}",
|
||||
}
|
||||
|
||||
result = await create_agent_and_generate_response(config)
|
||||
return result.choices[0]["message"]["content"]
|
||||
async with aiohttp.ClientSession() as session:
|
||||
async with session.post(url, json=payload, headers=headers, timeout=aiohttp.ClientTimeout(total=300)) as resp:
|
||||
if resp.status != 200:
|
||||
body = await resp.text()
|
||||
raise RuntimeError(f"API returned {resp.status}: {body}")
|
||||
data = await resp.json()
|
||||
|
||||
return data["choices"][0]["message"]["content"]
|
||||
|
||||
def _update_task_after_execution(self, task_id: str, tasks_file: Path):
|
||||
"""执行后更新 tasks.yaml"""
|
||||
|
||||
Loading…
Reference in New Issue
Block a user