|
|
@@ -1,15 +1,19 @@
|
|
|
/// <reference path="../env.d.ts" />
|
|
|
import { tool } from "@opencode-ai/plugin"
|
|
|
+import { homedir } from "node:os"
|
|
|
+import path from "node:path"
|
|
|
|
|
|
-const AGENTPAAS_URL = process.env.AGENTPAAS_URL ?? "http://127.0.0.1:8000"
|
|
|
+const AGENTPAAS_CONFIG_PATH = path.join(homedir(), ".agentpaas", "config.json")
|
|
|
+const DEFAULT_AGENTPAAS_URL = "http://127.0.0.1:8000"
|
|
|
|
|
|
type SSEEvent = { event: string; data: unknown }
|
|
|
+type AgentListItem = { id: string; name: string }
|
|
|
|
|
|
export default tool({
|
|
|
description: `把任务转发给远端 AgentPaaS 智能体执行,等待完成后返回最终答案。
|
|
|
|
|
|
参数说明:
|
|
|
-- agent_id:远端 AgentPaaS 智能体 ID(ag_ 开头),必填
|
|
|
+- agent_name:远端 AgentPaaS 智能体名称,工具会动态解析当前 ID,必填
|
|
|
- input:要执行的用户任务指令
|
|
|
- mode:iterate(多轮迭代,默认)/ chat(单轮)/ edit(单轮并指定目标子智能体)
|
|
|
- target_subagent:mode=edit 时的目标子智能体 ID
|
|
|
@@ -20,12 +24,9 @@ export default tool({
|
|
|
适用于调用远端科研流程智能体(如 research-67 系列)的场景。`,
|
|
|
args: {
|
|
|
input: tool.schema.string().describe("用户任务指令"),
|
|
|
- agent_id: tool.schema.string().describe("远端 AgentPaaS 智能体 ID(ag_ 开头)"),
|
|
|
+ agent_name: tool.schema.string().describe("远端 AgentPaaS 智能体名称"),
|
|
|
mode: tool.schema.enum(["iterate", "chat", "edit"]).default("iterate"),
|
|
|
- target_subagent: tool.schema
|
|
|
- .string()
|
|
|
- .optional()
|
|
|
- .describe("mode=edit 时的目标子智能 ID"),
|
|
|
+ target_subagent: tool.schema.string().optional().describe("mode=edit 时的目标子智能 ID"),
|
|
|
work_dir: tool.schema.string().optional().describe("覆盖工作目录(默认用当前项目目录)"),
|
|
|
thread_id: tool.schema.string().optional().describe("会话线程 ID,不传则新建"),
|
|
|
parameters: tool.schema
|
|
|
@@ -34,20 +35,22 @@ export default tool({
|
|
|
.describe("额外参数,原样传给远端智能体"),
|
|
|
},
|
|
|
async execute(args, context) {
|
|
|
- const apiKey = process.env.AGENTPAAS_API_KEY
|
|
|
- if (!apiKey) {
|
|
|
- throw new Error("AGENTPAAS_API_KEY 未设置:请先设置环境变量 AGENTPAAS_API_KEY(AgentPaaS Bearer Token)再启动 opencode")
|
|
|
- }
|
|
|
+ const config = await loadAgentPaaSConfig()
|
|
|
if (args.mode === "edit" && !args.target_subagent) {
|
|
|
throw new Error("mode=edit 时必须提供 target_subagent")
|
|
|
}
|
|
|
+ const agentID = await resolveAgentID(config, args.agent_name, context.abort)
|
|
|
const workDir = args.work_dir ?? context.worktree ?? context.directory
|
|
|
- const url = `${AGENTPAAS_URL}/api/v1/agents/${args.agent_id}/run/stream`
|
|
|
+ const url = `${config.url}/api/v1/agents/${encodeURIComponent(agentID)}/run/stream`
|
|
|
let resp: Response
|
|
|
try {
|
|
|
resp = await fetch(url, {
|
|
|
method: "POST",
|
|
|
- headers: { Authorization: `Bearer ${apiKey}`, "Content-Type": "application/json", Accept: "text/event-stream" },
|
|
|
+ headers: {
|
|
|
+ Authorization: `Bearer ${config.apiKey}`,
|
|
|
+ "Content-Type": "application/json",
|
|
|
+ Accept: "text/event-stream",
|
|
|
+ },
|
|
|
body: JSON.stringify({
|
|
|
input: args.input,
|
|
|
mode: args.mode,
|
|
|
@@ -68,13 +71,130 @@ export default tool({
|
|
|
if (!resp.body) throw new Error("AgentPaaS 响应没有 body")
|
|
|
|
|
|
const errors: string[] = []
|
|
|
+ const liveMetadata: Record<string, unknown> = {
|
|
|
+ agent_name: args.agent_name,
|
|
|
+ agent_id: agentID,
|
|
|
+ status: "running",
|
|
|
+ }
|
|
|
let done: Record<string, unknown> | null = null
|
|
|
let sawCancelled = false
|
|
|
+ let lastThinkUpdate = 0
|
|
|
+ const updateMetadata = (title: string, next: Record<string, unknown>) => {
|
|
|
+ Object.assign(liveMetadata, next)
|
|
|
+ context.metadata({ title, metadata: { ...liveMetadata } })
|
|
|
+ }
|
|
|
for await (const ev of parseSSE(resp.body)) {
|
|
|
- if (ev.event === "error") errors.push(extractMessage(ev.data))
|
|
|
- else if (ev.event === "cancelled") sawCancelled = true
|
|
|
- else if (ev.event === "done" && ev.data && typeof ev.data === "object") {
|
|
|
- done = ev.data as Record<string, unknown>
|
|
|
+ const data = asRecord(ev.data)
|
|
|
+ if (ev.event === "started") {
|
|
|
+ updateMetadata("Research67 正在运行", {
|
|
|
+ run_id: data?.run_id,
|
|
|
+ thread_id: data?.thread_id,
|
|
|
+ latest_event: ev.event,
|
|
|
+ })
|
|
|
+ } else if (ev.event === "stage_started") {
|
|
|
+ const stageName = getString(data?.stage_name) || getString(data?.stage_id) || "未知阶段"
|
|
|
+ updateMetadata(`Research67 正在进行 ${stageName}`, {
|
|
|
+ pipeline_id: data?.pipeline_id,
|
|
|
+ stage_id: data?.stage_id,
|
|
|
+ stage_name: stageName,
|
|
|
+ stage_status: data?.status,
|
|
|
+ attempt: data?.attempt,
|
|
|
+ latest_event: ev.event,
|
|
|
+ })
|
|
|
+ } else if (ev.event === "stage_passed") {
|
|
|
+ const stageName = getString(data?.stage_name) || getString(data?.stage_id) || "当前阶段"
|
|
|
+ updateMetadata(`Research67 已通过 ${stageName}`, {
|
|
|
+ pipeline_id: data?.pipeline_id,
|
|
|
+ stage_id: data?.stage_id,
|
|
|
+ stage_name: stageName,
|
|
|
+ stage_status: data?.status,
|
|
|
+ attempt: data?.attempt,
|
|
|
+ latest_event: ev.event,
|
|
|
+ })
|
|
|
+ } else if (ev.event === "stage_failed") {
|
|
|
+ const stageName = getString(data?.stage_name) || getString(data?.stage_id) || "当前阶段"
|
|
|
+ updateMetadata(`Research67 阶段失败:${stageName}`, {
|
|
|
+ pipeline_id: data?.pipeline_id,
|
|
|
+ stage_id: data?.stage_id,
|
|
|
+ stage_name: stageName,
|
|
|
+ stage_status: data?.status,
|
|
|
+ stage_error: data?.error,
|
|
|
+ attempt: data?.attempt,
|
|
|
+ latest_event: ev.event,
|
|
|
+ status: "failed",
|
|
|
+ })
|
|
|
+ } else if (ev.event === "pipeline_completed") {
|
|
|
+ updateMetadata("Research67 流水线已完成", {
|
|
|
+ pipeline_id: data?.pipeline_id,
|
|
|
+ pipeline_status: data?.status,
|
|
|
+ latest_event: ev.event,
|
|
|
+ status: "completed",
|
|
|
+ })
|
|
|
+ } else if (ev.event === "confirm_required") {
|
|
|
+ const runID = getString(data?.run_id) || getString(liveMetadata.run_id)
|
|
|
+ const toolName = getString(data?.tool) || "高风险工具"
|
|
|
+ const reason = getString(data?.reason)
|
|
|
+ const input = getString(data?.input)
|
|
|
+ if (!runID) throw new Error("AgentPaaS confirm_required 事件缺少 run_id")
|
|
|
+
|
|
|
+ updateMetadata(`Research67 正在等待确认:${toolName}`, {
|
|
|
+ latest_event: ev.event,
|
|
|
+ pending_confirmation: true,
|
|
|
+ confirmation_tool: toolName,
|
|
|
+ confirmation_reason: reason,
|
|
|
+ confirmation_input: input,
|
|
|
+ })
|
|
|
+ let approved = true
|
|
|
+ try {
|
|
|
+ await context.ask({
|
|
|
+ permission: "agentpaas_confirm",
|
|
|
+ patterns: [`${agentID}:${toolName}`],
|
|
|
+ always: [`${agentID}:${toolName}`],
|
|
|
+ metadata: { run_id: runID, agent_name: args.agent_name, tool: toolName, reason, input },
|
|
|
+ })
|
|
|
+ } catch {
|
|
|
+ approved = false
|
|
|
+ }
|
|
|
+ await resolveConfirmation(config, runID, approved, context.abort)
|
|
|
+ updateMetadata(`Research67 已${approved ? "批准" : "拒绝"} ${toolName}`, {
|
|
|
+ latest_event: "confirmation_resolved",
|
|
|
+ pending_confirmation: false,
|
|
|
+ confirmation_approved: approved,
|
|
|
+ })
|
|
|
+ } else if (ev.event === "tool_call") {
|
|
|
+ const toolName = getString(data?.tool) || "工具"
|
|
|
+ updateMetadata(withStage(`Research67 正在调用 ${toolName}`, liveMetadata), {
|
|
|
+ latest_event: ev.event,
|
|
|
+ latest_tool: toolName,
|
|
|
+ step: data?.step,
|
|
|
+ })
|
|
|
+ } else if (ev.event === "tool_result") {
|
|
|
+ const toolName = getString(data?.tool) || getString(liveMetadata.latest_tool) || "工具"
|
|
|
+ updateMetadata(withStage(`Research67 已完成 ${toolName}`, liveMetadata), {
|
|
|
+ latest_event: ev.event,
|
|
|
+ latest_tool: toolName,
|
|
|
+ step: data?.step,
|
|
|
+ })
|
|
|
+ } else if (ev.event === "think" || ev.event === "think_chunk") {
|
|
|
+ const now = Date.now()
|
|
|
+ if (ev.event === "think" || now - lastThinkUpdate >= 1000) {
|
|
|
+ lastThinkUpdate = now
|
|
|
+ updateMetadata(withStage("Research67 正在思考", liveMetadata), {
|
|
|
+ latest_event: ev.event,
|
|
|
+ step: data?.step,
|
|
|
+ })
|
|
|
+ }
|
|
|
+ } else if (ev.event === "answer") {
|
|
|
+ updateMetadata(withStage("Research67 正在整理结果", liveMetadata), {
|
|
|
+ latest_event: ev.event,
|
|
|
+ step: data?.step,
|
|
|
+ })
|
|
|
+ } else if (ev.event === "error") errors.push(extractMessage(ev.data))
|
|
|
+ else if (ev.event === "cancelled") {
|
|
|
+ sawCancelled = true
|
|
|
+ updateMetadata("Research67 已取消", { latest_event: ev.event, status: "cancelled" })
|
|
|
+ } else if (ev.event === "done" && data) {
|
|
|
+ done = data
|
|
|
break
|
|
|
}
|
|
|
}
|
|
|
@@ -84,19 +204,111 @@ export default tool({
|
|
|
}
|
|
|
const status = typeof done.status === "string" ? done.status : "unknown"
|
|
|
const metadata = {
|
|
|
- agent_id: args.agent_id, run_id: done.run_id, thread_id: done.thread_id, status,
|
|
|
- steps: done.steps, total_tokens: done.total_tokens, cost_usd: done.cost_usd, workspace_path: done.workspace_path,
|
|
|
+ agent_name: args.agent_name,
|
|
|
+ agent_id: agentID,
|
|
|
+ run_id: done.run_id,
|
|
|
+ thread_id: done.thread_id,
|
|
|
+ status,
|
|
|
+ steps: done.steps,
|
|
|
+ total_tokens: done.total_tokens,
|
|
|
+ cost_usd: done.cost_usd,
|
|
|
+ workspace_path: done.workspace_path,
|
|
|
}
|
|
|
- if (status === "completed") return { title: `AgentPaaS ${args.agent_id}`, output: String(done.output ?? ""), metadata }
|
|
|
- if (status === "cancelled") return { title: `AgentPaaS ${args.agent_id}`, output: "任务已被取消。", metadata }
|
|
|
+ if (status === "completed")
|
|
|
+ return { title: `AgentPaaS ${args.agent_name}`, output: String(done.output ?? ""), metadata }
|
|
|
+ if (status === "cancelled") return { title: `AgentPaaS ${args.agent_name}`, output: "任务已被取消。", metadata }
|
|
|
return {
|
|
|
- title: `AgentPaaS ${args.agent_id}(失败)`,
|
|
|
+ title: `AgentPaaS ${args.agent_name}(失败)`,
|
|
|
output: `[AgentPaaS 执行失败] ${errors.join("; ") || String(done.output ?? "") || "未知错误"}`,
|
|
|
metadata,
|
|
|
}
|
|
|
},
|
|
|
})
|
|
|
|
|
|
+async function resolveAgentID(config: { url: string; apiKey: string }, agentName: string, signal: AbortSignal) {
|
|
|
+ const name = agentName.trim()
|
|
|
+ if (!name) throw new Error("agent_name 不能为空")
|
|
|
+
|
|
|
+ let resp: Response
|
|
|
+ try {
|
|
|
+ resp = await fetch(`${config.url}/api/v1/agents`, {
|
|
|
+ headers: { Authorization: `Bearer ${config.apiKey}`, Accept: "application/json" },
|
|
|
+ signal,
|
|
|
+ })
|
|
|
+ } catch (err) {
|
|
|
+ throw new Error(`获取 AgentPaaS Agent 列表失败: ${err instanceof Error ? err.message : String(err)}`)
|
|
|
+ }
|
|
|
+ if (!resp.ok) {
|
|
|
+ const detail = await resp.text().catch(() => "")
|
|
|
+ throw new Error(`获取 AgentPaaS Agent 列表失败: HTTP ${resp.status}: ${detail.slice(0, 500)}`)
|
|
|
+ }
|
|
|
+
|
|
|
+ const payload: unknown = await resp.json().catch(() => null)
|
|
|
+ if (!payload || typeof payload !== "object" || !("agents" in payload) || !Array.isArray(payload.agents)) {
|
|
|
+ throw new Error("AgentPaaS Agent 列表响应格式无效")
|
|
|
+ }
|
|
|
+ const agents = (payload.agents as unknown[]).flatMap((item): AgentListItem[] => {
|
|
|
+ if (!item || typeof item !== "object" || !("id" in item) || !("name" in item)) return []
|
|
|
+ if (typeof item.id !== "string" || typeof item.name !== "string") return []
|
|
|
+ return [{ id: item.id, name: item.name }]
|
|
|
+ })
|
|
|
+ const matches = agents.filter((item) => item.name === name)
|
|
|
+ if (matches.length === 0) throw new Error(`AgentPaaS 中没有名为“${name}”的 active Agent`)
|
|
|
+ if (matches.length > 1)
|
|
|
+ throw new Error(`AgentPaaS 中存在 ${matches.length} 个名为“${name}”的 active Agent,无法确定调用目标`)
|
|
|
+ if (!matches[0].id.startsWith("ag_")) throw new Error(`AgentPaaS 为“${name}”返回了无效 ID: ${matches[0].id}`)
|
|
|
+ return matches[0].id
|
|
|
+}
|
|
|
+
|
|
|
+async function loadAgentPaaSConfig() {
|
|
|
+ const file = Bun.file(AGENTPAAS_CONFIG_PATH)
|
|
|
+ let saved: Record<string, unknown> = {}
|
|
|
+ if (await file.exists()) {
|
|
|
+ try {
|
|
|
+ const parsed: unknown = await file.json()
|
|
|
+ if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) throw new Error("根节点必须是 JSON 对象")
|
|
|
+ saved = parsed as Record<string, unknown>
|
|
|
+ } catch (err) {
|
|
|
+ throw new Error(
|
|
|
+ `AgentPaaS 配置文件无效:${AGENTPAAS_CONFIG_PATH}\n${err instanceof Error ? err.message : String(err)}`,
|
|
|
+ )
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ const savedUrl = typeof saved.server === "string" ? saved.server.trim() : ""
|
|
|
+ const savedApiKey = typeof saved.api_key === "string" ? saved.api_key.trim() : ""
|
|
|
+ const url = (process.env.AGENTPAAS_URL?.trim() || savedUrl || DEFAULT_AGENTPAAS_URL).replace(/\/+$/, "")
|
|
|
+ const apiKey = process.env.AGENTPAAS_API_KEY?.trim() || savedApiKey
|
|
|
+ if (apiKey) return { url, apiKey }
|
|
|
+
|
|
|
+ throw new Error(`AgentPaaS 尚未配置。请创建 ${AGENTPAAS_CONFIG_PATH}:
|
|
|
+{
|
|
|
+ "server": "${DEFAULT_AGENTPAAS_URL}",
|
|
|
+ "api_key": "你的 AgentPaaS API Key"
|
|
|
+}
|
|
|
+也可以通过 AGENTPAAS_URL 和 AGENTPAAS_API_KEY 环境变量临时覆盖。`)
|
|
|
+}
|
|
|
+
|
|
|
+async function resolveConfirmation(
|
|
|
+ config: { url: string; apiKey: string },
|
|
|
+ runID: string,
|
|
|
+ approved: boolean,
|
|
|
+ signal: AbortSignal,
|
|
|
+) {
|
|
|
+ const resp = await fetch(`${config.url}/api/v1/traces/${encodeURIComponent(runID)}/confirm`, {
|
|
|
+ method: "POST",
|
|
|
+ headers: {
|
|
|
+ Authorization: `Bearer ${config.apiKey}`,
|
|
|
+ "Content-Type": "application/json",
|
|
|
+ },
|
|
|
+ body: JSON.stringify({ approved }),
|
|
|
+ signal,
|
|
|
+ })
|
|
|
+ if (resp.ok) return
|
|
|
+ const detail = await resp.text().catch(() => "")
|
|
|
+ throw new Error(`AgentPaaS 确认请求失败: HTTP ${resp.status}: ${detail.slice(0, 500)}`)
|
|
|
+}
|
|
|
+
|
|
|
async function* parseSSE(body: ReadableStream<Uint8Array>): AsyncGenerator<SSEEvent> {
|
|
|
const reader = body.getReader()
|
|
|
const decoder = new TextDecoder()
|
|
|
@@ -114,7 +326,11 @@ async function* parseSSE(body: ReadableStream<Uint8Array>): AsyncGenerator<SSEEv
|
|
|
const dataRaw = parseSSEField(chunk, "data")
|
|
|
let data: unknown = dataRaw
|
|
|
if (dataRaw !== "") {
|
|
|
- try { data = JSON.parse(dataRaw) } catch { /* 非 JSON 当纯文本 */ }
|
|
|
+ try {
|
|
|
+ data = JSON.parse(dataRaw)
|
|
|
+ } catch {
|
|
|
+ /* 非 JSON 当纯文本 */
|
|
|
+ }
|
|
|
}
|
|
|
if (event) yield { event, data }
|
|
|
}
|
|
|
@@ -133,4 +349,20 @@ function parseSSEField(chunk: string, field: string): string {
|
|
|
function extractMessage(data: unknown): string {
|
|
|
if (data && typeof data === "object" && "message" in data) return String((data as Record<string, unknown>).message)
|
|
|
return String(data)
|
|
|
-}
|
|
|
+}
|
|
|
+
|
|
|
+function asRecord(value: unknown): Record<string, unknown> | undefined {
|
|
|
+ if (!value || typeof value !== "object" || Array.isArray(value)) return undefined
|
|
|
+ const result: Record<string, unknown> = {}
|
|
|
+ for (const [key, item] of Object.entries(value)) result[key] = item
|
|
|
+ return result
|
|
|
+}
|
|
|
+
|
|
|
+function getString(value: unknown): string {
|
|
|
+ return typeof value === "string" ? value : ""
|
|
|
+}
|
|
|
+
|
|
|
+function withStage(title: string, metadata: Record<string, unknown>): string {
|
|
|
+ const stageName = getString(metadata.stage_name)
|
|
|
+ return stageName ? `${title}(${stageName})` : title
|
|
|
+}
|