internal/opencode: runner (idle/hard timeout, process-group kill, resume-fallback) + verdict/json parsing
All checks were successful
build-test / build (push) Successful in 1m21s
All checks were successful
build-test / build (push) Successful in 1m21s
This commit is contained in:
16
go.mod
16
go.mod
@@ -1,3 +1,17 @@
|
|||||||
module github.com/kamelion/ratatoskr-go
|
module github.com/kamelion/ratatoskr-go
|
||||||
|
|
||||||
go 1.23
|
go 1.25.0
|
||||||
|
|
||||||
|
require modernc.org/sqlite v1.56.0
|
||||||
|
|
||||||
|
require (
|
||||||
|
github.com/dustin/go-humanize v1.0.1 // indirect
|
||||||
|
github.com/google/uuid v1.6.0 // indirect
|
||||||
|
github.com/mattn/go-isatty v0.0.24 // indirect
|
||||||
|
github.com/ncruces/go-strftime v1.0.0 // indirect
|
||||||
|
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect
|
||||||
|
golang.org/x/sys v0.47.0 // indirect
|
||||||
|
modernc.org/libc v1.74.4 // indirect
|
||||||
|
modernc.org/mathutil v1.7.1 // indirect
|
||||||
|
modernc.org/memory v1.11.0 // indirect
|
||||||
|
)
|
||||||
|
|||||||
50
go.sum
Normal file
50
go.sum
Normal file
@@ -0,0 +1,50 @@
|
|||||||
|
github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY=
|
||||||
|
github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto=
|
||||||
|
github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3 h1:LMLX+LgTNWpfvCBdFebv6EsYotImrt/Ppc5cXIriCSo=
|
||||||
|
github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3/go.mod h1:jl5iWTm0/hd5PjEYEOuwAJ57L/CibdZfrqZ5XA5GrCk=
|
||||||
|
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
|
||||||
|
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
|
||||||
|
github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k=
|
||||||
|
github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM=
|
||||||
|
github.com/mattn/go-isatty v0.0.24 h1:tGZZoVgT/KiqK1c8ocVLeDS8BSWMRd47J3Lbz7vsReI=
|
||||||
|
github.com/mattn/go-isatty v0.0.24/go.mod h1:nMCL3Zebbrt45jsMDgnfIwz6ydEQApk5oEI3HqDio6A=
|
||||||
|
github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w=
|
||||||
|
github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls=
|
||||||
|
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE=
|
||||||
|
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo=
|
||||||
|
golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ=
|
||||||
|
golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0=
|
||||||
|
golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM=
|
||||||
|
golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
|
||||||
|
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
|
||||||
|
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||||
|
golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q=
|
||||||
|
golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA=
|
||||||
|
modernc.org/cc/v4 v4.29.1 h1:MKgdCV3WykTSPqpVrnxdEDS0HEd2FHpKZDzxzU5LyeI=
|
||||||
|
modernc.org/cc/v4 v4.29.1/go.mod h1:OnovgIhbbMXMu1aISnJ0wvVD1KnW+cAUJkIrAWh+kVI=
|
||||||
|
modernc.org/ccgo/v4 v4.34.6 h1:sBgfIwyN0TQ9C5hwIeuqyeAKyMWnbvj2fvpF4L11uzU=
|
||||||
|
modernc.org/ccgo/v4 v4.34.6/go.mod h1:SZ8YcN9NG7XVsQYdm6jYBvi8PQP1qi+kqB6OhjqI3Fk=
|
||||||
|
modernc.org/fileutil v1.4.0 h1:j6ZzNTftVS054gi281TyLjHPp6CPHr2KCxEXjEbD6SM=
|
||||||
|
modernc.org/fileutil v1.4.0/go.mod h1:EqdKFDxiByqxLk8ozOxObDSfcVOv/54xDs/DUHdvCUU=
|
||||||
|
modernc.org/gc/v2 v2.6.5 h1:nyqdV8q46KvTpZlsw66kWqwXRHdjIlJOhG6kxiV/9xI=
|
||||||
|
modernc.org/gc/v2 v2.6.5/go.mod h1:YgIahr1ypgfe7chRuJi2gD7DBQiKSLMPgBQe9oIiito=
|
||||||
|
modernc.org/gc/v3 v3.1.4 h1:2g65LGVSmFQrXeITAw97x7hCRvZFcyE1uDP+7Vng7JI=
|
||||||
|
modernc.org/gc/v3 v3.1.4/go.mod h1:HFK/6AGESC7Ex+EZJhJ2Gni6cTaYpSMmU/cT9RmlfYY=
|
||||||
|
modernc.org/goabi0 v0.2.0 h1:HvEowk7LxcPd0eq6mVOAEMai46V+i7Jrj13t4AzuNks=
|
||||||
|
modernc.org/goabi0 v0.2.0/go.mod h1:CEFRnnJhKvWT1c1JTI3Avm+tgOWbkOu5oPA8eH8LnMI=
|
||||||
|
modernc.org/libc v1.74.4 h1:fX1Omw4o2/1C2iRkkIsrQTasJQldLhRmuPreXLoWs9k=
|
||||||
|
modernc.org/libc v1.74.4/go.mod h1:eeQAS9W3sZeKYMFubydxJpII9ybHWshk+7or7bLG9co=
|
||||||
|
modernc.org/mathutil v1.7.1 h1:GCZVGXdaN8gTqB1Mf/usp1Y/hSqgI2vAGGP4jZMCxOU=
|
||||||
|
modernc.org/mathutil v1.7.1/go.mod h1:4p5IwJITfppl0G4sUEDtCr4DthTaT47/N3aT6MhfgJg=
|
||||||
|
modernc.org/memory v1.11.0 h1:o4QC8aMQzmcwCK3t3Ux/ZHmwFPzE6hf2Y5LbkRs+hbI=
|
||||||
|
modernc.org/memory v1.11.0/go.mod h1:/JP4VbVC+K5sU2wZi9bHoq2MAkCnrt2r98UGeSK7Mjw=
|
||||||
|
modernc.org/opt v0.2.0 h1:tGyef5ApycA7FSEOMraay9SaTk5zmbx7Tu+cJs4QKZg=
|
||||||
|
modernc.org/opt v0.2.0/go.mod h1:03fq9lsNfvkYSfxrfUhZCWPk1lm4cq4N+Bh//bEtgns=
|
||||||
|
modernc.org/sortutil v1.2.1 h1:+xyoGf15mM3NMlPDnFqrteY07klSFxLElE2PVuWIJ7w=
|
||||||
|
modernc.org/sortutil v1.2.1/go.mod h1:7ZI3a3REbai7gzCLcotuw9AC4VZVpYMjDzETGsSMqJE=
|
||||||
|
modernc.org/sqlite v1.56.0 h1:/D8e2RfFqoy/Zc6PuC76U28zFwmI/sYx1Kjm4yEn9e0=
|
||||||
|
modernc.org/sqlite v1.56.0/go.mod h1:yCJ2cmAaIkHQ25oXWrF8H4O1lIfPYPR26yCEDj2P3pQ=
|
||||||
|
modernc.org/strutil v1.2.1 h1:UneZBkQA+DX2Rp35KcM69cSsNES9ly8mQWD71HKlOA0=
|
||||||
|
modernc.org/strutil v1.2.1/go.mod h1:EHkiggD70koQxjVdSBM3JKM7k6L0FbGE5eymy9i3B9A=
|
||||||
|
modernc.org/token v1.1.0 h1:Xl7Ap9dKaEs5kLoOQeQmPWevfnk/DM5qcLcYlA8ys6Y=
|
||||||
|
modernc.org/token v1.1.0/go.mod h1:UGzOrNV1mAFSEB63lOFHIpNRUVMvYTc6yu1SMY/XTDM=
|
||||||
86
internal/opencode/extract.go
Normal file
86
internal/opencode/extract.go
Normal file
@@ -0,0 +1,86 @@
|
|||||||
|
// Package opencode — запуск opencode-субагентов и разбор их вердиктов.
|
||||||
|
//
|
||||||
|
// Единственный компонент, умеющий запускать opencode run и парсить вывод.
|
||||||
|
// Контракты перенесены 1-в-1 из Python-версии (extract.py / opencode.py):
|
||||||
|
// - ExtractVerdict: последний text-парт из NDJSON-потока opencode run --format json
|
||||||
|
// - ExtractJSON: fenced ```json``` → первый {...}
|
||||||
|
// - Run/ResumeDev: запуск процесса с idle/hard timeout по opencode.db
|
||||||
|
package opencode
|
||||||
|
|
||||||
|
import (
|
||||||
|
"encoding/json"
|
||||||
|
"regexp"
|
||||||
|
"strings"
|
||||||
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
fenceRe = regexp.MustCompile("```(?:json)?\\s*([\\s\\S]*?)```")
|
||||||
|
jsonBlockRe = regexp.MustCompile("\\{[\\s\\S]*\\}")
|
||||||
|
sessionRe = regexp.MustCompile(`"session_id"\s*:\s*"([^"]+)"`)
|
||||||
|
)
|
||||||
|
|
||||||
|
// ExtractVerdict возвращает текст вердикта из NDJSON-потока opencode run --format json.
|
||||||
|
// Поток — NDJSON: последний парт {"type":"text","part":{"text":...}} и есть вердикт.
|
||||||
|
// Если JSON-партов нет (например, простой текст) — возвращает исходную строку.
|
||||||
|
func ExtractVerdict(out string) string {
|
||||||
|
out = strings.TrimSpace(out)
|
||||||
|
if out == "" {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
var last string
|
||||||
|
for _, line := range strings.Split(out, "\n") {
|
||||||
|
line = strings.TrimSpace(line)
|
||||||
|
if !strings.HasPrefix(line, "{") {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
var obj struct {
|
||||||
|
Type string `json:"type"`
|
||||||
|
Part struct {
|
||||||
|
Text string `json:"text"`
|
||||||
|
} `json:"part"`
|
||||||
|
}
|
||||||
|
if err := json.Unmarshal([]byte(line), &obj); err != nil {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if obj.Type == "text" && obj.Part.Text != "" {
|
||||||
|
last = obj.Part.Text
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if last != "" {
|
||||||
|
return strings.TrimSpace(last)
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
// ExtractJSON извлекает JSON из ответа модели: сначала fenced ```json```,
|
||||||
|
// затем первый {...}. Возвращает (nil, false), если JSON нет или невалиден.
|
||||||
|
// Это НЕ ошибка — вызывающий решает, как деградировать (класс O3 ErrParse).
|
||||||
|
func ExtractJSON(text string) (map[string]json.RawMessage, bool) {
|
||||||
|
if strings.TrimSpace(text) == "" {
|
||||||
|
return nil, false
|
||||||
|
}
|
||||||
|
for _, f := range fenceRe.FindAllStringSubmatch(text, -1) {
|
||||||
|
if len(f) > 1 {
|
||||||
|
var obj map[string]json.RawMessage
|
||||||
|
if err := json.Unmarshal([]byte(strings.TrimSpace(f[1])), &obj); err == nil && obj != nil {
|
||||||
|
return obj, true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if m := jsonBlockRe.FindString(text); m != "" {
|
||||||
|
var obj map[string]json.RawMessage
|
||||||
|
if err := json.Unmarshal([]byte(m), &obj); err == nil && obj != nil {
|
||||||
|
return obj, true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil, false
|
||||||
|
}
|
||||||
|
|
||||||
|
// SessionIDFromOutput извлекает session_id из текстового вывода opencode.
|
||||||
|
func SessionIDFromOutput(out string) (string, bool) {
|
||||||
|
m := sessionRe.FindStringSubmatch(out)
|
||||||
|
if len(m) > 1 && m[1] != "" {
|
||||||
|
return m[1], true
|
||||||
|
}
|
||||||
|
return "", false
|
||||||
|
}
|
||||||
115
internal/opencode/extract_test.go
Normal file
115
internal/opencode/extract_test.go
Normal file
@@ -0,0 +1,115 @@
|
|||||||
|
package opencode
|
||||||
|
|
||||||
|
import (
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestExtractVerdict(t *testing.T) {
|
||||||
|
cases := []struct {
|
||||||
|
name string
|
||||||
|
in string
|
||||||
|
want string
|
||||||
|
}{
|
||||||
|
{
|
||||||
|
name: "один text-парт NDJSON",
|
||||||
|
in: `{"type":"text","part":{"text":"hello"}}`,
|
||||||
|
want: "hello",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "несколько партов — берём последний text",
|
||||||
|
in: "{\"type\":\"step_start\",\"part\":{}}\n{\"type\":\"text\",\"part\":{\"text\":\"первый\"}}\n{\"type\":\"reasoning\",\"part\":{}}\n{\"type\":\"text\",\"part\":{\"text\":\"финал\"}}",
|
||||||
|
want: "финал",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "нет JSON — возвращаем строку как есть",
|
||||||
|
in: "простой текст без json",
|
||||||
|
want: "простой текст без json",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "пустая строка",
|
||||||
|
in: "",
|
||||||
|
want: "",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "text с пустым текстом пропускается",
|
||||||
|
in: "{\"type\":\"text\",\"part\":{\"text\":\"\"}}\n{\"type\":\"text\",\"part\":{\"text\":\"вал\"}}",
|
||||||
|
want: "вал",
|
||||||
|
},
|
||||||
|
}
|
||||||
|
for _, c := range cases {
|
||||||
|
t.Run(c.name, func(t *testing.T) {
|
||||||
|
if got := ExtractVerdict(c.in); got != c.want {
|
||||||
|
t.Errorf("ExtractVerdict(%q) = %q, want %q", c.in, got, c.want)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestExtractJSON(t *testing.T) {
|
||||||
|
cases := []struct {
|
||||||
|
name string
|
||||||
|
in string
|
||||||
|
want map[string]string
|
||||||
|
ok bool
|
||||||
|
}{
|
||||||
|
{
|
||||||
|
name: "fenced json",
|
||||||
|
in: "Вот ответ:\n```json\n{\"decision\":\"go\"}\n```",
|
||||||
|
want: map[string]string{"decision": "go"},
|
||||||
|
ok: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "fenced json без тега",
|
||||||
|
in: "```\n{\"a\":1}\n```",
|
||||||
|
want: map[string]string{"a": "1"},
|
||||||
|
ok: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "голый json-объект в тексте",
|
||||||
|
in: "Решение: {\"phase\":\"propose\"}",
|
||||||
|
want: map[string]string{"phase": "propose"},
|
||||||
|
ok: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "нет json",
|
||||||
|
in: "просто текст",
|
||||||
|
ok: false,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "невалидный json в fence",
|
||||||
|
in: "```json\n{no:json}\n```",
|
||||||
|
ok: false,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "пустая строка",
|
||||||
|
in: " ",
|
||||||
|
ok: false,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
for _, c := range cases {
|
||||||
|
t.Run(c.name, func(t *testing.T) {
|
||||||
|
got, ok := ExtractJSON(c.in)
|
||||||
|
if ok != c.ok {
|
||||||
|
t.Fatalf("ExtractJSON(%q) ok = %v, want %v", c.in, ok, c.ok)
|
||||||
|
}
|
||||||
|
if !ok {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
for k, v := range c.want {
|
||||||
|
if !strings.Contains(string(got[k]), v) {
|
||||||
|
t.Errorf("ExtractJSON(%q)[%s] = %s, want to contain %s", c.in, k, got[k], v)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSessionIDFromOutput(t *testing.T) {
|
||||||
|
if s, ok := SessionIDFromOutput(`{"session_id":"abc123"}`); !ok || s != "abc123" {
|
||||||
|
t.Fatalf("got %q %v", s, ok)
|
||||||
|
}
|
||||||
|
if _, ok := SessionIDFromOutput("no session here"); ok {
|
||||||
|
t.Fatal("expected no match")
|
||||||
|
}
|
||||||
|
}
|
||||||
23
internal/opencode/pgid_linux.go
Normal file
23
internal/opencode/pgid_linux.go
Normal file
@@ -0,0 +1,23 @@
|
|||||||
|
//go:build linux
|
||||||
|
|
||||||
|
package opencode
|
||||||
|
|
||||||
|
import (
|
||||||
|
"os/exec"
|
||||||
|
"syscall"
|
||||||
|
)
|
||||||
|
|
||||||
|
func sysProcAttr(proc *exec.Cmd) {
|
||||||
|
proc.SysProcAttr = &syscall.SysProcAttr{Setpgid: true}
|
||||||
|
}
|
||||||
|
|
||||||
|
// killProcGroup убивает всю process-group по лидеру pid (SIGKILL дочерним и
|
||||||
|
// SIGTERM лидеру). Игнорирует ошибки: weakest-effort teardown.
|
||||||
|
func killProcGroup(pid int) {
|
||||||
|
pgid, err := syscall.Getpgid(pid)
|
||||||
|
if err != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
_ = syscall.Kill(-pgid, syscall.SIGKILL)
|
||||||
|
_ = syscall.Kill(pid, syscall.SIGKILL)
|
||||||
|
}
|
||||||
9
internal/opencode/pgid_other.go
Normal file
9
internal/opencode/pgid_other.go
Normal file
@@ -0,0 +1,9 @@
|
|||||||
|
//go:build !linux
|
||||||
|
|
||||||
|
package opencode
|
||||||
|
|
||||||
|
import "os/exec"
|
||||||
|
|
||||||
|
func sysProcAttr(_ *exec.Cmd) {}
|
||||||
|
|
||||||
|
func killProcGroup(pid int) {}
|
||||||
252
internal/opencode/runner.go
Normal file
252
internal/opencode/runner.go
Normal file
@@ -0,0 +1,252 @@
|
|||||||
|
package opencode
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bufio"
|
||||||
|
"context"
|
||||||
|
"database/sql"
|
||||||
|
"fmt"
|
||||||
|
"io"
|
||||||
|
"os"
|
||||||
|
"os/exec"
|
||||||
|
"strings"
|
||||||
|
"sync"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
_ "modernc.org/sqlite" // чисто-Go драйвер, без CGO → один статический бинарь
|
||||||
|
)
|
||||||
|
|
||||||
|
// Result — результат запуска opencode run. rc=-1 означает «убит по таймауту»
|
||||||
|
// (idle/hard): вызывающий НЕ должен ронять задачу, а обязан закоммитить/запушить
|
||||||
|
// готовую работу и отправить на ревью (класс O2 Timeout — результат, не ошибка).
|
||||||
|
type Result struct {
|
||||||
|
RC int
|
||||||
|
Stdout string
|
||||||
|
SessionID string
|
||||||
|
}
|
||||||
|
|
||||||
|
// Runner — конфигурация запуска opencode-субагентов.
|
||||||
|
type Runner struct {
|
||||||
|
Bin string // путь к opencode (по умолчанию "opencode")
|
||||||
|
DBPath string // путь к opencode.db (idle-детекция активности)
|
||||||
|
Config string // путь к opencode.json (OPENCODE_CONFIG)
|
||||||
|
IdleTimeout time.Duration // нет новых сообщений в БД → зависание
|
||||||
|
HardTimeout time.Duration // общий лимит на запуск
|
||||||
|
PollInterval time.Duration
|
||||||
|
|
||||||
|
// Заменяемые для тестов:
|
||||||
|
Stdout io.Writer // диагностика (лог), по умолчанию os.Stderr
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *Runner) defaults() {
|
||||||
|
if r.Bin == "" {
|
||||||
|
r.Bin = "opencode"
|
||||||
|
}
|
||||||
|
if r.IdleTimeout == 0 {
|
||||||
|
r.IdleTimeout = 2 * time.Minute
|
||||||
|
}
|
||||||
|
if r.HardTimeout == 0 {
|
||||||
|
r.HardTimeout = 20 * time.Minute
|
||||||
|
}
|
||||||
|
if r.PollInterval == 0 {
|
||||||
|
r.PollInterval = 2 * time.Second
|
||||||
|
}
|
||||||
|
if r.Stdout == nil {
|
||||||
|
r.Stdout = os.Stderr
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *Runner) logf(format string, args ...any) {
|
||||||
|
fmt.Fprintf(r.Stdout, format+"\n", args...)
|
||||||
|
}
|
||||||
|
|
||||||
|
// maxDirMsgTS — максимальный time_updated (мс) по всем сообщениям сессий этого
|
||||||
|
// worktree: сигнал «модель/субагенты ещё активны». nil-nil если БД нет/пуста.
|
||||||
|
func (r *Runner) maxDirMsgTS(ctx context.Context, worktree string) (int64, bool) {
|
||||||
|
if r.DBPath == "" {
|
||||||
|
return 0, false
|
||||||
|
}
|
||||||
|
db, err := sql.Open("sqlite", "file:"+r.DBPath+"?mode=ro")
|
||||||
|
if err != nil {
|
||||||
|
return 0, false
|
||||||
|
}
|
||||||
|
defer db.Close()
|
||||||
|
var ts sql.NullInt64
|
||||||
|
err = db.QueryRowContext(ctx,
|
||||||
|
"SELECT MAX(m.time_updated) FROM message m JOIN session s ON s.id = m.session_id WHERE s.directory = ?",
|
||||||
|
worktree).Scan(&ts)
|
||||||
|
if err != nil || !ts.Valid {
|
||||||
|
return 0, false
|
||||||
|
}
|
||||||
|
return ts.Int64, true
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *Runner) latestSession(ctx context.Context, worktree, agent string) (string, bool) {
|
||||||
|
if r.DBPath == "" {
|
||||||
|
return "", false
|
||||||
|
}
|
||||||
|
db, err := sql.Open("sqlite", "file:"+r.DBPath+"?mode=ro")
|
||||||
|
if err != nil {
|
||||||
|
return "", false
|
||||||
|
}
|
||||||
|
defer db.Close()
|
||||||
|
q := "SELECT id FROM session WHERE directory = ?"
|
||||||
|
args := []any{worktree}
|
||||||
|
if agent != "" {
|
||||||
|
q += " AND agent = ?"
|
||||||
|
args = append(args, agent)
|
||||||
|
}
|
||||||
|
q += " ORDER BY time_created DESC LIMIT 1"
|
||||||
|
var id string
|
||||||
|
if err := db.QueryRowContext(ctx, q, args...).Scan(&id); err != nil {
|
||||||
|
return "", false
|
||||||
|
}
|
||||||
|
return id, id != ""
|
||||||
|
}
|
||||||
|
|
||||||
|
// Run запускает opencode run. Возвращает *Result (rc, stdout, session_id).
|
||||||
|
// Ошибка — только класс O1 ErrSpawn (не смог запустить бинарь). Таймауты
|
||||||
|
// дают rc=-1 в Result, а не error (класс O2).
|
||||||
|
func (r *Runner) Run(ctx context.Context, prompt, cwd, agent, sessionID string) (*Result, error) {
|
||||||
|
r.defaults()
|
||||||
|
cmd := []string{r.Bin, "run", "--agent", agent, "--format", "json", "--dir", cwd}
|
||||||
|
if sessionID != "" {
|
||||||
|
cmd = append(cmd, "--session", sessionID)
|
||||||
|
}
|
||||||
|
cmd = append(cmd, prompt)
|
||||||
|
|
||||||
|
env := append(os.Environ(),
|
||||||
|
"OPENCODE_DISABLE_AUTOUPDATE=1",
|
||||||
|
"OPENCODE_DISABLE_MODELS_FETCH=1")
|
||||||
|
if r.Config != "" {
|
||||||
|
env = append(env, "OPENCODE_CONFIG="+r.Config)
|
||||||
|
}
|
||||||
|
|
||||||
|
proc := exec.CommandContext(ctx, cmd[0], cmd[1:]...)
|
||||||
|
proc.Env = env
|
||||||
|
proc.Dir = cwd
|
||||||
|
// Убиваем всю process-group, чтобы дочерние процессы (sleep и т.п.) тоже
|
||||||
|
// умерли и закрыли унаследованные stdout-fd (иначе <-done виснет).
|
||||||
|
setpgid(proc)
|
||||||
|
stdout, err := proc.StdoutPipe()
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("opencode: stdout pipe: %w", err)
|
||||||
|
}
|
||||||
|
proc.Stderr = proc.Stdout
|
||||||
|
if err := proc.Start(); err != nil {
|
||||||
|
return nil, fmt.Errorf("opencode: start %v: %w", cmd[0], err)
|
||||||
|
}
|
||||||
|
|
||||||
|
var buf []string
|
||||||
|
var mu sync.Mutex
|
||||||
|
done := make(chan struct{})
|
||||||
|
go func() {
|
||||||
|
defer close(done)
|
||||||
|
sc := bufio.NewScanner(stdout)
|
||||||
|
for sc.Scan() {
|
||||||
|
mu.Lock()
|
||||||
|
buf = append(buf, sc.Text())
|
||||||
|
mu.Unlock()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
baseline, _ := r.maxDirMsgTS(ctx, cwd)
|
||||||
|
lastProgress := time.Now()
|
||||||
|
launch := time.Now()
|
||||||
|
killed := false
|
||||||
|
|
||||||
|
pollLoop:
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-done:
|
||||||
|
// процесс завершился (pipe EOF) — выходим, берём exit code
|
||||||
|
break pollLoop
|
||||||
|
case <-ctx.Done():
|
||||||
|
killGroup(proc)
|
||||||
|
killed = true
|
||||||
|
break pollLoop
|
||||||
|
default:
|
||||||
|
}
|
||||||
|
if proc.ProcessState != nil && proc.ProcessState.Exited() {
|
||||||
|
break pollLoop
|
||||||
|
}
|
||||||
|
now := time.Now()
|
||||||
|
ts, ok := r.maxDirMsgTS(ctx, cwd)
|
||||||
|
if ok && ts > baseline {
|
||||||
|
lastProgress = now
|
||||||
|
}
|
||||||
|
if now.Sub(lastProgress) > r.IdleTimeout {
|
||||||
|
r.logf("opencode(%s) idle %.0fs (нет новых сообщений) — kill", agent, r.IdleTimeout.Seconds())
|
||||||
|
killGroup(proc)
|
||||||
|
killed = true
|
||||||
|
break pollLoop
|
||||||
|
}
|
||||||
|
if now.Sub(launch) > r.HardTimeout {
|
||||||
|
r.logf("opencode(%s) hard timeout %.0fs — kill", agent, r.HardTimeout.Seconds())
|
||||||
|
killGroup(proc)
|
||||||
|
killed = true
|
||||||
|
break pollLoop
|
||||||
|
}
|
||||||
|
time.Sleep(r.PollInterval)
|
||||||
|
}
|
||||||
|
|
||||||
|
<-done
|
||||||
|
procErr := proc.Wait()
|
||||||
|
rc := proc.ProcessState.ExitCode()
|
||||||
|
if rc < 0 {
|
||||||
|
rc = 1
|
||||||
|
}
|
||||||
|
if killed {
|
||||||
|
rc = -1
|
||||||
|
}
|
||||||
|
_ = procErr
|
||||||
|
|
||||||
|
mu.Lock()
|
||||||
|
out := strings.Join(buf, "\n")
|
||||||
|
mu.Unlock()
|
||||||
|
|
||||||
|
sid := sessionID
|
||||||
|
if s, ok := SessionIDFromOutput(out); ok {
|
||||||
|
sid = s
|
||||||
|
}
|
||||||
|
if rc == -1 && sid == "" {
|
||||||
|
if s, ok := r.latestSession(ctx, cwd, agent); ok {
|
||||||
|
sid = s
|
||||||
|
}
|
||||||
|
}
|
||||||
|
r.logf("opencode(%s) rc=%d", agent, rc)
|
||||||
|
return &Result{RC: rc, Stdout: out, SessionID: sid}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ResumeDev — запуск dev-агента с resume-fallback. Если resume (sessionID)
|
||||||
|
// падает с rc!=0 (напр. сессия потеряна) — повторяем ОДИН раз свежей сессией
|
||||||
|
// в том же worktree. rc=-1 (kill по таймауту) НЕ триггерит fallback.
|
||||||
|
// Возвращает (result, timedOut).
|
||||||
|
func (r *Runner) ResumeDev(ctx context.Context, prompt, cwd, sessionID string) (*Result, bool) {
|
||||||
|
res, err := r.Run(ctx, prompt, cwd, "dev", sessionID)
|
||||||
|
if err != nil {
|
||||||
|
// spawn-ошибку не ретраим fallback'ом — она повторится
|
||||||
|
return res, false
|
||||||
|
}
|
||||||
|
if res.RC != 0 && res.RC != -1 && sessionID != "" {
|
||||||
|
r.logf("dev resume rc=%d — запускаю заново без --session (worktree сохраняю)", res.RC)
|
||||||
|
res, _ = r.Run(ctx, prompt+resumeFallbackNote, cwd, "dev", "")
|
||||||
|
}
|
||||||
|
return res, res.RC == -1
|
||||||
|
}
|
||||||
|
|
||||||
|
const resumeFallbackNote = "\n\n(Возобновление сессии не удалось; продолжи с учётом уже сделанных изменений в worktree.)"
|
||||||
|
|
||||||
|
// --- process-group helpers (Linux) ---
|
||||||
|
// Ставим процесс в собственную process-group, чтобы killGroup мог убить и
|
||||||
|
// дочерние процессы (иначе они держат унаследованные stdout-fd и <-done виснет).
|
||||||
|
|
||||||
|
func setpgid(proc *exec.Cmd) {
|
||||||
|
sysProcAttr(proc)
|
||||||
|
}
|
||||||
|
|
||||||
|
func killGroup(proc *exec.Cmd) {
|
||||||
|
if proc.Process != nil {
|
||||||
|
killProcGroup(proc.Process.Pid)
|
||||||
|
}
|
||||||
|
_ = proc.Process.Kill()
|
||||||
|
}
|
||||||
131
internal/opencode/runner_test.go
Normal file
131
internal/opencode/runner_test.go
Normal file
@@ -0,0 +1,131 @@
|
|||||||
|
package opencode
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// fakeOpenCode создаёт shell-скрипт, имитирующий opencode run:
|
||||||
|
// $FAKE_MODE=ok -> мгновенный успех, печатает NDJSON c session_id
|
||||||
|
// $FAKE_MODE=slow-> спит долго (для idle/hard timeout)
|
||||||
|
// $FAKE_MODE=fail-> exit 7 (resume-fallback)
|
||||||
|
func fakeOpenCode(t *testing.T, workdir string) string {
|
||||||
|
t.Helper()
|
||||||
|
bin := filepath.Join(workdir, "opencode")
|
||||||
|
script := `#!/bin/sh
|
||||||
|
mode="${FAKE_MODE:-ok}"
|
||||||
|
case "$mode" in
|
||||||
|
ok)
|
||||||
|
echo '{"type":"text","part":{"text":"done"}}'
|
||||||
|
echo '{"session_id":"sess-123"}'
|
||||||
|
exit 0
|
||||||
|
;;
|
||||||
|
slow)
|
||||||
|
sleep 30
|
||||||
|
;;
|
||||||
|
fail)
|
||||||
|
echo '{"type":"text","part":{"text":"boom"}}'
|
||||||
|
exit 7
|
||||||
|
;;
|
||||||
|
esac
|
||||||
|
`
|
||||||
|
if err := os.WriteFile(bin, []byte(script), 0o755); err != nil {
|
||||||
|
t.Fatalf("write fake opencode: %v", err)
|
||||||
|
}
|
||||||
|
return bin
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRun_Success(t *testing.T) {
|
||||||
|
dir := t.TempDir()
|
||||||
|
bin := fakeOpenCode(t, dir)
|
||||||
|
t.Setenv("FAKE_MODE", "ok")
|
||||||
|
|
||||||
|
r := &Runner{Bin: bin, PollInterval: 20 * time.Millisecond}
|
||||||
|
res, err := r.Run(context.Background(), "task", dir, "dev", "")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Run err: %v", err)
|
||||||
|
}
|
||||||
|
if res.RC != 0 {
|
||||||
|
t.Errorf("RC = %d, want 0", res.RC)
|
||||||
|
}
|
||||||
|
if res.SessionID != "sess-123" {
|
||||||
|
t.Errorf("SessionID = %q, want sess-123", res.SessionID)
|
||||||
|
}
|
||||||
|
if !contains(res.Stdout, "done") {
|
||||||
|
t.Errorf("Stdout = %q, want to contain done", res.Stdout)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRun_IdleTimeout(t *testing.T) {
|
||||||
|
dir := t.TempDir()
|
||||||
|
bin := fakeOpenCode(t, dir)
|
||||||
|
t.Setenv("FAKE_MODE", "slow")
|
||||||
|
|
||||||
|
r := &Runner{Bin: bin, IdleTimeout: 50 * time.Millisecond,
|
||||||
|
PollInterval: 10 * time.Millisecond}
|
||||||
|
res, err := r.Run(context.Background(), "task", dir, "dev", "")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Run err: %v", err)
|
||||||
|
}
|
||||||
|
if res.RC != -1 {
|
||||||
|
t.Errorf("RC = %d, want -1 (timeout kill)", res.RC)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRun_ContextCancel(t *testing.T) {
|
||||||
|
dir := t.TempDir()
|
||||||
|
bin := fakeOpenCode(t, dir)
|
||||||
|
t.Setenv("FAKE_MODE", "slow")
|
||||||
|
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
r := &Runner{Bin: bin, HardTimeout: time.Minute,
|
||||||
|
PollInterval: 10 * time.Millisecond}
|
||||||
|
done := make(chan *Result, 1)
|
||||||
|
errCh := make(chan error, 1)
|
||||||
|
go func() {
|
||||||
|
res, err := r.Run(ctx, "task", dir, "dev", "")
|
||||||
|
done <- res
|
||||||
|
errCh <- err
|
||||||
|
}()
|
||||||
|
time.Sleep(30 * time.Millisecond)
|
||||||
|
cancel()
|
||||||
|
res := <-done
|
||||||
|
if err := <-errCh; err != nil {
|
||||||
|
t.Fatalf("Run err: %v", err)
|
||||||
|
}
|
||||||
|
if res.RC != -1 {
|
||||||
|
t.Errorf("RC = %d, want -1", res.RC)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestResumeDev_Fallback(t *testing.T) {
|
||||||
|
dir := t.TempDir()
|
||||||
|
bin := fakeOpenCode(t, dir)
|
||||||
|
t.Setenv("FAKE_MODE", "fail")
|
||||||
|
|
||||||
|
r := &Runner{Bin: bin, PollInterval: 20 * time.Millisecond}
|
||||||
|
res, timedOut := r.ResumeDev(context.Background(), "task", dir, "lost-session")
|
||||||
|
if timedOut {
|
||||||
|
t.Error("timedOut = true, want false")
|
||||||
|
}
|
||||||
|
// fake fail всегда exit 7, fallback тоже 7 — проверяем что RC от fallback-вызова
|
||||||
|
if res.RC != 7 {
|
||||||
|
t.Errorf("RC = %d, want 7 (fallback повтор с тем же кодом)", res.RC)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func contains(s, sub string) bool {
|
||||||
|
return len(s) >= len(sub) && (s == sub || len(s) > 0 && indexOf(s, sub) >= 0)
|
||||||
|
}
|
||||||
|
|
||||||
|
func indexOf(s, sub string) int {
|
||||||
|
for i := 0; i+len(sub) <= len(s); i++ {
|
||||||
|
if s[i:i+len(sub)] == sub {
|
||||||
|
return i
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return -1
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user