557 lines
22 KiB
Go
557 lines
22 KiB
Go
package agent
|
||
|
||
import (
|
||
"bytes"
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"io"
|
||
"log"
|
||
"os"
|
||
"os/exec"
|
||
"path/filepath"
|
||
"regexp"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/docker/docker/api/types/container"
|
||
"github.com/docker/docker/api/types/mount"
|
||
)
|
||
|
||
// 协作镜像名(构建方式见各自目录的 RUN.md)。
|
||
const (
|
||
coopImagePigo = "pigo-coop" // pigo 协作镜像(支持 openai + anthropic 协议)
|
||
coopImage = "pi-coop" // pi 协作镜像(支持 openai + anthropic 协议)
|
||
coopImageClaude = "claude-coop" // Claude Code 协作镜像(仅支持 anthropic 协议)
|
||
)
|
||
|
||
// 协作引擎到镜像 / 容器名前缀 / 展示标签 / 运行时目录的映射。
|
||
// 三种引擎共享同一套环境变量契约(MODEL/BASE_URL/API_KEY/PROTOCOL/TASK 等),
|
||
// supervisor 各自适配,控制层只需选择镜像并做协议约束 + baseURL 规范化。
|
||
// runtimeDir 是 local 模式下 supervisor.sh / prompts / extensions 所在的相对目录。
|
||
var coopEngines = map[string]struct {
|
||
image, namePrefix, label, runtimeDir string
|
||
}{
|
||
"pigo": {coopImagePigo, "pigo-coop-", "pigo", "pigo/coop"},
|
||
"pi": {coopImage, "pi-coop-", "pi", "pi-coop"},
|
||
"claude": {coopImageClaude, "claude-coop-", "claude code", "claude code"},
|
||
}
|
||
|
||
// codePattern 从任务描述中提取题目编号(如 "c-06" / "a-05" / "e1-01")。
|
||
var codePattern = regexp.MustCompile(`[A-Za-z0-9]{1,8}-\d{1,4}`)
|
||
|
||
// runCoop 异步启动 pi 协作任务:
|
||
// 校验参数与镜像后立即返回,容器创建 / 运行 / 等待 / 清理全部放到后台 goroutine。
|
||
// 任务结束后由 CoopManager 通知 Agent,主 agent 会自动收到一条
|
||
// "【协作任务完成通知】"消息并汇报结果,无需在前端长时间等待。
|
||
func (t *Toolset) runCoop(ctx context.Context, sessionID string, params map[string]any) ToolResult {
|
||
if t.coop == nil {
|
||
return ToolResult{Success: false, Output: "协作任务管理器未初始化"}
|
||
}
|
||
task, _ := params["task"].(string)
|
||
task = strings.TrimSpace(task)
|
||
if task == "" {
|
||
return ToolResult{Success: false, Output: "task 不能为空"}
|
||
}
|
||
|
||
// 完成协议:约束协作 agent 主动收尾,避免"已提交成功却没写 DONE 被强杀"、
|
||
// "解不出却空耗到超时"两类问题。supervisor 只在黑板根目录出现 DONE 标记时正常退出。
|
||
completionProtocol := `
|
||
|
||
【完成协议(务必严格遵守,决定协作能否正常收尾)】
|
||
- 找到 flag 并提交成功(平台响应 correct=true)后:用 blackboard 工具 action=done,在黑板根目录创建 DONE 标记,内容写入完成摘要(含 flag 值、提交响应、解题路径)。supervisor 检测到 DONE 即正常结束(exit 0);不创建 DONE 会空耗到超时被强杀,协作被视为失败。
|
||
- 若经充分尝试后确认本轮无法解出(目标不可达 / 无漏洞 / 试错过多):同样用 action=done 创建 DONE,内容开头写明「未解出」与已尝试内容,让调度方及时关闭靶机并切换下一题,不要空耗到超时。
|
||
- 已通关题目不要重复提交:平台对已通关题目的后续提交统一返回 correct:false(而非 duplicate),属正常现象、不影响已得分数,不要误判为失败。`
|
||
task += completionProtocol
|
||
|
||
// 从任务描述解析题目编号(如 "c-06"),供完成通知携带,便于控制层关靶机/切题
|
||
challengeCode := ""
|
||
if m := codePattern.FindString(task); m != "" {
|
||
challengeCode = m
|
||
}
|
||
roundMax := intParam(params, "round_max")
|
||
if roundMax <= 0 {
|
||
// 单 agent 默认 1 轮:一次运行完成全部工作,未完成则失败并重新下发
|
||
roundMax = 1
|
||
}
|
||
if roundMax > 30 {
|
||
roundMax = 30
|
||
}
|
||
timeoutSec := intParam(params, "timeout")
|
||
if timeoutSec <= 0 {
|
||
// 默认 1800s:600s/900s 对需要写脚本+多步探测的渗透/解题任务偏紧,
|
||
// 实测多因单轮超时(exit 143)导致协作失败。
|
||
timeoutSec = 1800
|
||
}
|
||
|
||
if t.apiCfg == nil {
|
||
return ToolResult{Success: false, Output: "缺少 LLM API 配置,无法注入模型配置"}
|
||
}
|
||
|
||
// 解析协作引擎:engine=pi(默认)/ pigo / claude。
|
||
// 未显式指定时使用设置页配置的默认引擎(apiCfg.Engine())。
|
||
// 三种镜像共享同一套环境变量契约(MODEL/BASE_URL/API_KEY/PROTOCOL/TASK 等),
|
||
// supervisor 各自适配,控制层只需选择镜像并做协议约束 + baseURL 规范化。
|
||
engine, _ := params["engine"].(string)
|
||
engine = strings.TrimSpace(strings.ToLower(engine))
|
||
if engine == "" {
|
||
engine = t.apiCfg.Engine()
|
||
}
|
||
eng, ok := coopEngines[engine]
|
||
if !ok {
|
||
return ToolResult{Success: false, Output: "不支持的 engine \"" + engine + "\",可选值:pigo | pi(默认)| claude"}
|
||
}
|
||
|
||
model := t.apiCfg.Model()
|
||
apiKey := t.apiCfg.APIKey()
|
||
provider := t.apiCfg.Provider()
|
||
// coopBaseURL 根据 engine 规范化 baseURL:
|
||
// pigo 的 anthropicCompatDriver 只追加 /messages,需补 /v1;
|
||
// pi/claude 用官方 SDK(自带 /v1/messages),需剥离 /v1 避免双重路径。
|
||
baseURL := coopBaseURL(provider, t.apiCfg.BaseURL(), engine)
|
||
if model == "" || apiKey == "" || baseURL == "" {
|
||
return ToolResult{Success: false, Output: "LLM API 配置不完整(model / base_url / api_key 缺一不可),请先在设置页配置"}
|
||
}
|
||
|
||
// Claude Code 仅支持 Anthropic 协议端点(ANTHROPIC_BASE_URL),
|
||
// 若配置为 openai 协议则直接拒绝,避免容器启动后才报错。
|
||
if engine == "claude" && provider != ProviderAnthropic {
|
||
return ToolResult{Success: false, Output: "claude-coop 仅支持 Anthropic 协议端点,当前 provider 为 " + provider +
|
||
"。请改用 engine=pi 或 engine=pigo,或将 API 配置切换为 anthropic 协议(如 DeepSeek 的 /anthropic 端点)。"}
|
||
}
|
||
|
||
env := []string{
|
||
"MODEL=" + model,
|
||
"BASE_URL=" + baseURL,
|
||
"API_KEY=" + apiKey,
|
||
"PROTOCOL=" + provider,
|
||
"TASK=" + task,
|
||
fmt.Sprintf("ROUND_MAX=%d", roundMax),
|
||
fmt.Sprintf("TIMEOUT=%d", timeoutSec),
|
||
}
|
||
|
||
// 生成任务 ID(也用作容器名与默认黑板子目录)
|
||
taskID := newID()
|
||
|
||
// 黑板目录:用户显式指定 blackboard 时使用指定路径;
|
||
// 未指定时为该任务分配独立的会话工作区子目录(data/workspaces/<sessionID>/blackboard/<taskID>),
|
||
// 保证同一会话发起的多个 coop 任务互相隔离、互不污染。
|
||
blackboardDir := ""
|
||
if dirParam, _ := params["blackboard"].(string); strings.TrimSpace(dirParam) != "" {
|
||
dirParam = strings.TrimSpace(dirParam)
|
||
abs, err := filepath.Abs(dirParam)
|
||
if err != nil {
|
||
return ToolResult{Success: false, Output: "blackboard 路径无效: " + err.Error()}
|
||
}
|
||
// 安全约束:blackboard 会被 worker 进程/容器读写,必须限制在项目根目录内,
|
||
// 否则 LLM 可通过指定任意主机目录(如 ~/.ssh)让 worker 读写敏感文件。
|
||
if !pathWithin(abs, t.Workspace) {
|
||
return ToolResult{Success: false, Output: "blackboard 路径超出项目根目录范围,已拒绝: " + abs}
|
||
}
|
||
if err := os.MkdirAll(abs, 0o755); err != nil {
|
||
return ToolResult{Success: false, Output: "创建 blackboard 目录失败: " + err.Error()}
|
||
}
|
||
blackboardDir = abs
|
||
} else {
|
||
dir := filepath.Join(t.workspaceFor(sessionID), "blackboard", taskID)
|
||
if err := os.MkdirAll(dir, 0o755); err != nil {
|
||
return ToolResult{Success: false, Output: "创建 blackboard 目录失败: " + err.Error()}
|
||
}
|
||
blackboardDir = dir
|
||
}
|
||
|
||
record := &CoopTask{
|
||
ID: taskID,
|
||
SessionID: sessionID,
|
||
Status: "running",
|
||
Blackboard: blackboardDir,
|
||
RoundMax: roundMax,
|
||
CreatedAt: time.Now(),
|
||
ChallengeCode: challengeCode,
|
||
}
|
||
t.coop.Register(record)
|
||
|
||
mode := t.apiCfg.CoopMode()
|
||
log.Printf("[coop] 协作任务启动 task=%s session=%s mode=%s engine=%s challenge=%s round_max=%d timeout=%ds blackboard=%s",
|
||
record.ID, sessionID, mode, engine, challengeCode, roundMax, timeoutSec, blackboardDir)
|
||
|
||
overallSec := timeoutSec*roundMax + 300
|
||
if mode == CoopModeLocal {
|
||
// 本地子进程模式:无需 Docker,直接 exec supervisor.sh
|
||
supervisor, promptsDir, extensionsDir, err := findCoopRuntime(engine)
|
||
if err != nil {
|
||
t.coop.Complete(t.failedTask(record, err))
|
||
return ToolResult{Success: false, Output: err.Error()}
|
||
}
|
||
go t.runCoopLocalAsync(record, env, blackboardDir, supervisor, promptsDir, extensionsDir, roundMax, timeoutSec, eng.label)
|
||
} else {
|
||
// Docker 容器模式:现有逻辑
|
||
var mounts []mount.Mount
|
||
cli, err := t.dockerClient()
|
||
if err != nil {
|
||
t.coop.Complete(t.failedTask(record, err))
|
||
return ToolResult{Success: false, Output: err.Error()}
|
||
}
|
||
defer cli.Close()
|
||
|
||
// 判断目标 daemon 系统类型:远程 Linux daemon(如 WSL)的 bind 挂载
|
||
// 只认容器侧路径,Windows 路径需转换为 /mnt/<盘符>/... 形式。
|
||
info, infoErr := cli.Info(ctx)
|
||
linuxDaemon := infoErr == nil && info.OSType == "linux"
|
||
|
||
// 检查镜像是否存在(同步快速失败,避免后台任务因镜像缺失白跑)
|
||
checkCtx, cancel := context.WithTimeout(ctx, t.Timeout)
|
||
defer cancel()
|
||
if _, err := cli.ImageInspect(checkCtx, eng.image); err != nil {
|
||
t.coop.Complete(t.failedTask(record, fmt.Errorf("未找到镜像 %s", eng.image)))
|
||
return ToolResult{Success: false, Output: "未找到镜像 " + eng.image + "。请先在仓库根构建:\n" +
|
||
coopBuildHint(engine) + "\n(详见对应目录的 RUN.md)"}
|
||
}
|
||
|
||
// 容器内以非 root 用户(uid=1000 agent)运行,主机目录必须对任何用户可写,
|
||
// 否则 supervisor 初始化(mkdir/写 task.md/复制 AGENTS.md)会失败。
|
||
_ = os.Chmod(blackboardDir, 0o777)
|
||
source := blackboardDir
|
||
if linuxDaemon {
|
||
source = wslBindPath(blackboardDir)
|
||
}
|
||
mounts = append(mounts, mount.Mount{
|
||
Type: mount.TypeBind,
|
||
Source: source,
|
||
Target: "/blackboard",
|
||
})
|
||
|
||
// 容器创建 / 启动 / 等待 / 清理放到后台,并使用独立上下文,
|
||
// 避免阻塞当前 SSE 流(此前同步等待最长可达 timeout×round_max+300 秒)。
|
||
go t.runCoopAsync(record, env, mounts, roundMax, timeoutSec, eng.image, eng.namePrefix)
|
||
}
|
||
|
||
return ToolResult{Success: true, Output: fmt.Sprintf(
|
||
"协作任务已在后台启动,任务 ID:%s。\n运行方式:%s 协作单 Agent(%s 模式),最多 %d 轮,整体上限约 %d 分钟。\n"+
|
||
"你无需在此等待,可以继续处理其他请求;任务完成后系统会自动通知你并汇报结果。",
|
||
record.ID, eng.label, mode, roundMax, overallSec/60)}
|
||
}
|
||
|
||
// runCoopAsync 在后台完成协作容器的创建、启动、等待、日志收集与清理。
|
||
// 结束(成功 / 失败 / 超时)后通过 coop.Complete 通知 Agent 唤醒主 agent。
|
||
// image / namePrefix 由 runCoop 根据 engine 选择(pi-coop 或 claude-coop)。
|
||
func (t *Toolset) runCoopAsync(record *CoopTask, env []string, mounts []mount.Mount, roundMax, timeoutSec int, image, namePrefix string) {
|
||
// 使用独立后台上下文,避免随 SSE 请求断开而中断协作任务
|
||
ctx := context.Background()
|
||
|
||
cli, err := t.dockerClient()
|
||
if err != nil {
|
||
t.coop.Complete(t.failedTask(record, err))
|
||
return
|
||
}
|
||
defer cli.Close()
|
||
|
||
name := namePrefix + record.ID[:8]
|
||
created, err := cli.ContainerCreate(ctx, &container.Config{
|
||
Image: image,
|
||
Env: env,
|
||
}, &container.HostConfig{
|
||
Mounts: mounts,
|
||
}, nil, nil, name)
|
||
if err != nil {
|
||
t.coop.Complete(t.failedTask(record, fmt.Errorf("创建容器失败: %w", err)))
|
||
return
|
||
}
|
||
record.ContainerID = created.ID
|
||
|
||
// 运行结束后无论如何清理容器(等价 docker run --rm)
|
||
cleanup := func() {
|
||
cleanupCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||
defer cancel()
|
||
_ = cli.ContainerRemove(cleanupCtx, created.ID, container.RemoveOptions{Force: true})
|
||
}
|
||
|
||
if err := cli.ContainerStart(ctx, created.ID, container.StartOptions{}); err != nil {
|
||
cleanup()
|
||
t.coop.Complete(t.failedTask(record, fmt.Errorf("启动容器失败: %w", err)))
|
||
return
|
||
}
|
||
|
||
// 等待容器退出;整体超时上限 = 轮次 × 单轮超时 + 缓冲
|
||
overall := time.Duration(timeoutSec*roundMax+300) * time.Second
|
||
waitCtx, waitCancel := context.WithTimeout(ctx, overall)
|
||
defer waitCancel()
|
||
waitCh, errCh := cli.ContainerWait(waitCtx, created.ID, container.WaitConditionNotRunning)
|
||
|
||
exitCode := -1
|
||
select {
|
||
case res := <-waitCh:
|
||
exitCode = int(res.StatusCode)
|
||
case err := <-errCh:
|
||
cleanup()
|
||
t.coop.Complete(t.failedTask(record, fmt.Errorf("等待容器退出失败: %w", err)))
|
||
return
|
||
case <-waitCtx.Done():
|
||
cleanup()
|
||
t.coop.Complete(t.failedTask(record, fmt.Errorf(
|
||
"协作运行超时(整体上限 %d 秒,约 %.0f 分钟),已强制清理容器。可增大 round_max / timeout 或缩小任务规模后重试。",
|
||
int(overall.Seconds()), overall.Minutes())))
|
||
return
|
||
}
|
||
|
||
// 读取容器日志(supervisor 输出 + DONE 总结)
|
||
logCtx, logCancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||
defer logCancel()
|
||
result := ""
|
||
if logs, err := cli.ContainerLogs(logCtx, created.ID, container.LogsOptions{
|
||
ShowStdout: true,
|
||
ShowStderr: true,
|
||
}); err == nil {
|
||
if raw, readErr := io.ReadAll(logs); readErr == nil {
|
||
result = demuxDockerLogs(raw)
|
||
_ = logs.Close()
|
||
}
|
||
}
|
||
|
||
result = strings.TrimSpace(result)
|
||
switch exitCode {
|
||
case 0:
|
||
result += "\n[协作完成,退出码 0]"
|
||
case 1:
|
||
result += "\n[达到最大轮次仍未完成,退出码 1;可增大 round_max 后重试]"
|
||
default:
|
||
result += fmt.Sprintf("\n[容器退出码 %d,请检查上方日志]", exitCode)
|
||
}
|
||
if len(mounts) > 0 {
|
||
result += "\n黑板产物已保留在: " + mounts[0].Source
|
||
}
|
||
cleanup()
|
||
|
||
record.ExitCode = exitCode
|
||
record.Output = t.truncate(result)
|
||
|
||
// 解析 supervisor 生成的结构化结果 result.json(如有):
|
||
// 完成通知与汇报轮优先使用 Result 字段,容器 stdout 仅作兜底。
|
||
status := "unknown"
|
||
if data, rerr := os.ReadFile(filepath.Join(record.Blackboard, "result.json")); rerr == nil {
|
||
var res CoopResult
|
||
if json.Unmarshal(data, &res) == nil {
|
||
record.Result = &res
|
||
status = res.Status
|
||
}
|
||
}
|
||
log.Printf("[coop] 协作任务结束 task=%s session=%s exit_code=%d status=%s", record.ID, record.SessionID, exitCode, status)
|
||
t.coop.Complete(record)
|
||
}
|
||
|
||
// failedTask 生成一个失败任务记录,供 runCoopAsync 失败路径统一上报。
|
||
func (t *Toolset) failedTask(record *CoopTask, err error) *CoopTask {
|
||
record.Error = err.Error()
|
||
record.Output = "协作任务失败: " + err.Error()
|
||
return record
|
||
}
|
||
|
||
// coopBuildHint 返回指定 engine 对应镜像的构建命令提示。
|
||
func coopBuildHint(engine string) string {
|
||
switch engine {
|
||
case "pigo":
|
||
// pigo-coop 构建上下文是 pigo/ 目录,且需先交叉编译 pigo 二进制
|
||
return "cd pigo && GOOS=linux CGO_ENABLED=0 go build -trimpath -ldflags=\"-s -w\" -o coop/tmp/pigo-linux-amd64 ./cmd/pigo && docker build -f coop/Dockerfile -t pigo-coop ."
|
||
case "claude":
|
||
return "docker build -f \"claude code/Dockerfile\" -t claude-coop \"claude code/\""
|
||
default:
|
||
return "docker build -f pi-coop/Dockerfile -t pi-coop pi-coop/"
|
||
}
|
||
}
|
||
|
||
// coopBaseURL 把当前 Agent 配置的 Base URL 规范化为协作容器期望的"基础地址"。
|
||
// 不同 engine 的模型接入层对 baseURL 的路径处理不同,需按 engine 区分:
|
||
// - pigo:anthropicCompatDriver 只追加 /messages(不含 /v1),因此 anthropic
|
||
// 端点需补 /v1(最终 /v1/messages);openai 端点需保留 /v1。
|
||
// - pi:provider 使用官方 SDK(@anthropic-ai/sdk / openai),SDK 自行追加完整
|
||
// 路径(anthropic: /v1/messages;openai: /chat/completions),不能再补 /v1,
|
||
// 否则产生 /v1/v1/messages 双重路径导致 404。
|
||
// - claude:同 pi,用 @anthropic-ai/sdk,不能补 /v1。
|
||
//
|
||
// 对用户误填的完整路径后缀(/v1/messages、/chat/completions 等)统一剥离,
|
||
// 再按 engine + provider 决定是否补 /v1。
|
||
func coopBaseURL(provider, baseURL, engine string) string {
|
||
base := strings.TrimRight(strings.TrimSpace(baseURL), "/")
|
||
// 先剥离用户误填的完整路径后缀,统一回退到"基础地址"
|
||
switch {
|
||
case strings.HasSuffix(base, "/v1/chat/completions"):
|
||
base = strings.TrimSuffix(base, "/chat/completions")
|
||
case strings.HasSuffix(base, "/chat/completions"):
|
||
base = strings.TrimSuffix(base, "/chat/completions")
|
||
case strings.HasSuffix(base, "/v1/messages"):
|
||
base = strings.TrimSuffix(base, "/v1/messages")
|
||
case strings.HasSuffix(base, "/messages"):
|
||
base = strings.TrimSuffix(base, "/messages")
|
||
}
|
||
|
||
if provider != ProviderAnthropic {
|
||
return base
|
||
}
|
||
|
||
// Anthropic 协议端点处理
|
||
switch engine {
|
||
case "pigo":
|
||
// pigo 的 anthropicCompatDriver 只追加 /messages,需补 /v1。
|
||
// 若用户已填 /v1 结尾则保持,否则补上。
|
||
if strings.HasSuffix(base, "/v1") {
|
||
return base
|
||
}
|
||
return base + "/v1"
|
||
default:
|
||
// pi / claude 用 @anthropic-ai/sdk,SDK 自带 /v1/messages,
|
||
// 不能补 /v1;若用户已填 /v1 则剥离(SDK 会补回完整的 /v1/messages)。
|
||
if strings.HasSuffix(base, "/v1") {
|
||
return strings.TrimSuffix(base, "/v1")
|
||
}
|
||
return base
|
||
}
|
||
}
|
||
|
||
// wslBindPath 把 Windows 绝对路径转换为 WSL 挂载路径(E:\path → /mnt/e/path),
|
||
// 供运行在 WSL 内的 Linux Docker daemon 做 bind 挂载。
|
||
func wslBindPath(path string) string {
|
||
if len(path) < 2 || path[1] != ':' {
|
||
return path
|
||
}
|
||
drive := strings.ToLower(path[:1])
|
||
rest := strings.ReplaceAll(path[2:], "\\", "/")
|
||
return "/mnt/" + drive + rest
|
||
}
|
||
|
||
// findCoopRuntime 查找指定 engine 的 supervisor.sh / prompts / extensions 目录。
|
||
// 查找顺序:环境变量 COOP_DIR > 相对于可执行文件 > 相对于工作目录。
|
||
// 返回的路径均为绝对路径,供本地子进程模式(runCoopLocalAsync)使用。
|
||
func findCoopRuntime(engine string) (supervisor, prompts, extensions string, err error) {
|
||
eng, ok := coopEngines[engine]
|
||
if !ok {
|
||
return "", "", "", fmt.Errorf("不支持的 engine: %s", engine)
|
||
}
|
||
dirName := eng.runtimeDir
|
||
|
||
// 候选基目录列表:COOP_DIR 环境变量 > 可执行文件同级/上级 > 当前工作目录
|
||
var candidates []string
|
||
if envDir := os.Getenv("COOP_DIR"); envDir != "" {
|
||
candidates = append(candidates, filepath.Join(envDir, dirName))
|
||
}
|
||
if exe, exeErr := os.Executable(); exeErr == nil {
|
||
exeDir := filepath.Dir(exe)
|
||
candidates = append(candidates, filepath.Join(exeDir, dirName))
|
||
candidates = append(candidates, filepath.Join(exeDir, "..", dirName))
|
||
}
|
||
if wd, wdErr := os.Getwd(); wdErr == nil {
|
||
candidates = append(candidates, filepath.Join(wd, dirName))
|
||
}
|
||
|
||
for _, base := range candidates {
|
||
sp := filepath.Join(base, "supervisor.sh")
|
||
if st, statErr := os.Stat(sp); statErr == nil && !st.IsDir() {
|
||
pp := filepath.Join(base, "prompts")
|
||
ep := filepath.Join(base, "extensions")
|
||
// prompts / extensions 可选:缺失时传空串,supervisor 用内置默认
|
||
return sp, pp, ep, nil
|
||
}
|
||
}
|
||
return "", "", "", fmt.Errorf(
|
||
"本地模式未找到 %s 的 supervisor.sh,已查找目录: %v\n"+
|
||
"请确保 %s 目录存在且包含 supervisor.sh,或设置 COOP_DIR 环境变量指向包含该目录的父目录",
|
||
engine, candidates, dirName)
|
||
}
|
||
|
||
// runCoopLocalAsync 在后台以本地子进程方式运行 supervisor.sh 完成 worker 任务。
|
||
// 与 runCoopAsync(Docker 模式)对应:无需 Docker daemon,直接 exec supervisor.sh,
|
||
// 通过 BLACKBOARD/PROMPTS/EXTENSIONS 环境变量指向本地路径。
|
||
// 结束(成功 / 失败 / 超时)后通过 coop.Complete 通知 Agent。
|
||
func (t *Toolset) runCoopLocalAsync(record *CoopTask, env []string, blackboardDir, supervisor, promptsDir, extensionsDir string, roundMax, timeoutSec int, label string) {
|
||
// 使用独立后台上下文,避免随 SSE 请求断开而中断协作任务
|
||
ctx := context.Background()
|
||
overall := time.Duration(timeoutSec*roundMax+300) * time.Second
|
||
runCtx, runCancel := context.WithTimeout(ctx, overall)
|
||
defer runCancel()
|
||
|
||
// Windows 上 bash 通常是 WSL bash,不认反斜杠路径(E:\foo → E:foo 被吞)。
|
||
// 需把传给 bash 的路径转为 /mnt/<盘符>/... 格式;Go 侧文件操作仍用原始路径。
|
||
toBashPath := func(p string) string {
|
||
if len(p) >= 2 && p[1] == ':' {
|
||
return wslBindPath(p)
|
||
}
|
||
return p
|
||
}
|
||
bashSupervisor := toBashPath(supervisor)
|
||
bashBlackboard := toBashPath(blackboardDir)
|
||
bashPrompts := toBashPath(promptsDir)
|
||
bashExtensions := toBashPath(extensionsDir)
|
||
|
||
// 构建子进程环境:继承父进程环境(PATH 等)+ 注入协作环境变量
|
||
procEnv := os.Environ()
|
||
procEnv = append(procEnv, env...)
|
||
procEnv = append(procEnv, "BLACKBOARD="+bashBlackboard)
|
||
if bashPrompts != "" {
|
||
procEnv = append(procEnv, "PROMPTS="+bashPrompts)
|
||
}
|
||
if bashExtensions != "" {
|
||
procEnv = append(procEnv, "EXTENSIONS="+bashExtensions)
|
||
}
|
||
|
||
cmd := exec.CommandContext(runCtx, "bash", bashSupervisor)
|
||
cmd.Env = procEnv
|
||
// stdout+stderr 合并捕获(supervisor 的日志输出)
|
||
var buf bytes.Buffer
|
||
cmd.Stdout = &buf
|
||
cmd.Stderr = &buf
|
||
|
||
log.Printf("[coop] 本地协作进程启动 task=%s pid=pending blackboard=%s supervisor=%s",
|
||
record.ID, bashBlackboard, bashSupervisor)
|
||
|
||
if err := cmd.Start(); err != nil {
|
||
t.coop.Complete(t.failedTask(record, fmt.Errorf("启动 supervisor 失败: %w", err)))
|
||
return
|
||
}
|
||
record.ProcessID = cmd.Process.Pid
|
||
log.Printf("[coop] 本地协作进程已启动 task=%s pid=%d", record.ID, record.ProcessID)
|
||
|
||
// 等待进程退出(exec.CommandContext 在 runCtx 超时时自动发送 SIGKILL)
|
||
waitErr := cmd.Wait()
|
||
exitCode := 0
|
||
if waitErr != nil {
|
||
if exitErr, ok := waitErr.(*exec.ExitError); ok {
|
||
exitCode = exitErr.ExitCode()
|
||
} else {
|
||
exitCode = -1
|
||
}
|
||
}
|
||
|
||
// 判断是否因超时被杀
|
||
timedOut := runCtx.Err() == context.DeadlineExceeded
|
||
|
||
result := strings.TrimSpace(buf.String())
|
||
switch {
|
||
case timedOut:
|
||
result += fmt.Sprintf("\n[协作运行超时(整体上限 %d 秒,约 %.0f 分钟),已强制终止进程]",
|
||
int(overall.Seconds()), overall.Minutes())
|
||
exitCode = 124
|
||
case exitCode == 0:
|
||
result += "\n[协作完成,退出码 0]"
|
||
case exitCode == 1:
|
||
result += "\n[达到最大轮次仍未完成,退出码 1;可增大 round_max 后重试]"
|
||
default:
|
||
result += fmt.Sprintf("\n[进程退出码 %d,请检查上方日志]", exitCode)
|
||
}
|
||
result += "\n黑板产物已保留在: " + blackboardDir
|
||
|
||
record.ExitCode = exitCode
|
||
record.Output = t.truncate(result)
|
||
|
||
// 解析 supervisor 生成的结构化结果 result.json(与 docker 模式一致)
|
||
status := "unknown"
|
||
if data, rerr := os.ReadFile(filepath.Join(blackboardDir, "result.json")); rerr == nil {
|
||
var res CoopResult
|
||
if json.Unmarshal(data, &res) == nil {
|
||
record.Result = &res
|
||
status = res.Status
|
||
}
|
||
}
|
||
log.Printf("[coop] 本地协作任务结束 task=%s pid=%d exit_code=%d status=%s",
|
||
record.ID, record.ProcessID, exitCode, status)
|
||
t.coop.Complete(record)
|
||
}
|