feat: opencode через HTTP API — пул serve-серверов вместо spawn/NDJSON
Runner теперь ходит к постоянным serve по HTTP API (v1.17+, /api):
- клиент Client (create/send/wait/abort/messages/verdict)
- Pool: по одному serve на каталог, ленивый подъём, root-сервер в worktree,
выделение портов, ReleaseTask при завершении задачи
- Run: CreateSession('ratatoskr-<агент>') -> Send -> поллинг Verdict из
text-частей assistant-сообщений; idle/hard таймауты дают RC=-1
- вердикт извлекается из последнего assistant text-парта (плоский text)
- тесты: unit на фейковом HTTP-сервере; e2e эмулирует serve через httptest,
агент определяется по title сессии
This commit is contained in:
@@ -1,53 +1,42 @@
|
||||
package opencode
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
_ "modernc.org/sqlite" // чисто-Go драйвер, без CGO → один статический бинарь
|
||||
)
|
||||
|
||||
// Result — результат запуска opencode run. rc=-1 означает «убит по таймауту»
|
||||
// (idle/hard): вызывающий НЕ должен ронять задачу, а обязан закоммитить/запушить
|
||||
// готовую работу и отправить на ревью (класс O2 Timeout — результат, не ошибка).
|
||||
// Result — результат запуска opencode-субагента через HTTP API. rc=-1 означает
|
||||
// «оборван по таймауту/контексту» (idle/hard): вызывающий НЕ должен ронять
|
||||
// задачу, а обязан закоммитить/запушить готовую работу и отправить на ревью
|
||||
// (класс O2 Timeout — результат, не ошибка).
|
||||
type Result struct {
|
||||
RC int
|
||||
Stdout string
|
||||
SessionID string
|
||||
}
|
||||
|
||||
// Runner — конфигурация запуска opencode-субагентов.
|
||||
// Runner — запуск opencode-субагентов через HTTP API serve.
|
||||
//
|
||||
// Полный переход на API: Runner ходит к opencode serve через Pool→Client
|
||||
// (нет spawn-модели, нет NDJSON). Агент идёт в сервер пула для своего каталога
|
||||
// (в нём запущен serve → он его project).
|
||||
type Runner struct {
|
||||
Bin string // путь к opencode (по умолчанию "opencode")
|
||||
DBPath string // путь к opencode.db (idle-детекция активности)
|
||||
Config string // путь к opencode.json (OPENCODE_CONFIG)
|
||||
ConfigDir string // путь к каталогу с агентами (OPENCODE_CONFIG_DIR)
|
||||
IdleTimeout time.Duration // нет активных live-строк в стриме И сообщений в БД → завис
|
||||
HardTimeout time.Duration // общий лимит на запуск
|
||||
Pool *Pool // пул serve-серверов (обязательный)
|
||||
IdleTimeout time.Duration
|
||||
HardTimeout time.Duration
|
||||
PollInterval time.Duration
|
||||
|
||||
// AttachURL — постоянный opencode serve (режим --attach). Непустое значение
|
||||
// переключает команду на `run --attach <url> ...`: тёплые модели/MCP, без
|
||||
// холодного старта на каждый вызов. Пусто — историческая spawn-модель.
|
||||
AttachURL string
|
||||
|
||||
// Заменяемые для тестов:
|
||||
// Заменяемый для тестов:
|
||||
Stdout io.Writer // диагностика (лог), по умолчанию os.Stderr
|
||||
}
|
||||
|
||||
func (r *Runner) defaults() {
|
||||
if r.Bin == "" {
|
||||
r.Bin = "opencode"
|
||||
}
|
||||
if r.IdleTimeout == 0 {
|
||||
r.IdleTimeout = 5 * time.Minute
|
||||
}
|
||||
@@ -66,236 +55,157 @@ 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).
|
||||
// Run запускает opencode-субагента через HTTP API: создаёт/продолжает сессию
|
||||
// в сервере пула для каталога cwd, отправляет промпт, ждёт вердикт.
|
||||
//
|
||||
// Возвращает *Result (rc, stdout=вердикт, session_id). Ошибка — только класс
|
||||
// O1 ErrRun (не смог обратиться к серверу/сессии). Таймауты дают 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"}
|
||||
if r.AttachURL != "" {
|
||||
cmd = append(cmd, "--attach", r.AttachURL)
|
||||
}
|
||||
cmd = append(cmd, "--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)
|
||||
}
|
||||
if r.ConfigDir != "" {
|
||||
env = append(env, "OPENCODE_CONFIG_DIR="+r.ConfigDir)
|
||||
if r.Pool == nil {
|
||||
return nil, fmt.Errorf("opencode: Pool не задан (API-режим обязателен)")
|
||||
}
|
||||
|
||||
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()
|
||||
srv, err := r.Pool.Ensure(ctx, cwd)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("opencode: stdout pipe: %w", err)
|
||||
return nil, err
|
||||
}
|
||||
proc.Stderr = proc.Stdout
|
||||
if err := proc.Start(); err != nil {
|
||||
return nil, fmt.Errorf("opencode: start %v: %w", cmd[0], err)
|
||||
c := &Client{BaseURL: srv.Addr(), Password: srv.Password}
|
||||
|
||||
// Сессия: заданная (resume) или новая.
|
||||
sid := sessionID
|
||||
if sid == "" {
|
||||
sid, err = c.CreateSession(ctx, "ratatoskr-"+agent)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("opencode: create session: %w", err)
|
||||
}
|
||||
r.logf("opencode(%s) session=%s на %s", agent, sid, srv.Addr())
|
||||
}
|
||||
|
||||
var buf []string
|
||||
// Отправляем промпт (неблокирующий — сервер начинает выполнение).
|
||||
if err := c.Send(ctx, sid, prompt); err != nil {
|
||||
return nil, fmt.Errorf("opencode: prompt: %w", err)
|
||||
}
|
||||
|
||||
return r.awaitVerdict(ctx, c, sid, agent)
|
||||
}
|
||||
|
||||
// awaitVerdict поллит сообщения сессии, пока не появится готовый text-вердикт
|
||||
// от assistant, либо не истечёт idle/hard таймаут (тогда Abort + rc=-1).
|
||||
func (r *Runner) awaitVerdict(ctx context.Context, c *Client, sid, agent string) (*Result, error) {
|
||||
// Прогресс = сумма text-частей во всех assistant-сообщениях сессии. Рост
|
||||
// сбрасывает idle-таймер (LLM стримит = жив).
|
||||
var mu sync.Mutex
|
||||
done := make(chan struct{})
|
||||
// liveSeq — кол-во распознанных live-строк (text/tool/agent/reasoning) в
|
||||
// NDJSON-потоке. Инкрементится из goroutine чтения; поллинг сравнивает,
|
||||
// чтобы сбросить idle-таймер «пока LLM стримит» (а не только по БД).
|
||||
var liveSeq atomic.Uint64
|
||||
prevLive := liveSeq.Load()
|
||||
// Живое наблюдение сессии (если задано через WithLive в контексте).
|
||||
liveReg, liveTask := liveFromContext(ctx)
|
||||
if liveReg != nil && liveTask != 0 {
|
||||
liveReg.Start(liveTask, agent)
|
||||
defer liveReg.Finish(liveTask)
|
||||
}
|
||||
go func() {
|
||||
defer close(done)
|
||||
sc := bufio.NewScanner(stdout)
|
||||
// NDJSON opencode пишет каждый объект одной строкой; большой text-парт с
|
||||
// вердиктом легко превышает дефолтный лимит Scanner в 64КБ → ErrTooLong и
|
||||
// потеря всего потока после первой строки. Поднимаем до 64МБ.
|
||||
sc.Buffer(make([]byte, 64*1024), 64*1024*1024)
|
||||
for sc.Scan() {
|
||||
line := sc.Text()
|
||||
mu.Lock()
|
||||
buf = append(buf, line)
|
||||
mu.Unlock()
|
||||
if st := parseLiveStep(line); st != nil {
|
||||
// «пульс» LLM: что-то стримится/вызывается — сбрасываем idle
|
||||
liveSeq.Add(1)
|
||||
if liveReg != nil {
|
||||
liveReg.Observe(liveTask, *st)
|
||||
}
|
||||
}
|
||||
}
|
||||
scanErr := sc.Err()
|
||||
if scanErr != nil {
|
||||
r.logf("opencode(%s) scan err: %v", agent, scanErr)
|
||||
}
|
||||
}()
|
||||
|
||||
baseline, _ := r.maxDirMsgTS(ctx, cwd)
|
||||
lastCount := -1
|
||||
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 ctx.Err() != nil {
|
||||
if err := c.Abort(ctx, sid); err != nil {
|
||||
r.logf("opencode(%s) abort (ctx): %v", agent, err)
|
||||
}
|
||||
r.logf("opencode(%s) ctx cancelled — обрыв (rc=-1)", agent)
|
||||
return &Result{RC: -1, Stdout: "", SessionID: sid}, nil
|
||||
}
|
||||
if proc.ProcessState != nil && proc.ProcessState.Exited() {
|
||||
break pollLoop
|
||||
|
||||
count, _ := c.textCount(ctx, sid)
|
||||
mu.Lock()
|
||||
if count != lastCount {
|
||||
lastProgress = time.Now()
|
||||
lastCount = count
|
||||
}
|
||||
mu.Unlock()
|
||||
|
||||
// Пробуем вердикт (дешёвый GET). Если готов — выходим.
|
||||
vd, vErr := c.Verdict(ctx, sid)
|
||||
if vErr == nil && vd != "" {
|
||||
r.logf("opencode(%s) вердикт готов (%d байт)", agent, len(vd))
|
||||
return &Result{RC: 0, Stdout: vd, SessionID: sid}, nil
|
||||
}
|
||||
// Сервер недоступен — фатально (не таймаут). Но если это следствие
|
||||
// отмены контекста (cancel прилетел прямо во время запроса) — это обрыв
|
||||
// по контексту, а не недоступность сервера: вернём RC=-1 на следующей
|
||||
// итерации (проверка ctx.Err() вверху) и не подменяем класс ошибки.
|
||||
var ce *ClientErr
|
||||
if errors.As(vErr, &ce) && ce.Op == "connect" && ctx.Err() == nil {
|
||||
return nil, fmt.Errorf("opencode: %w", vErr)
|
||||
}
|
||||
|
||||
now := time.Now()
|
||||
// «Пульс» LLM: если с прошлого поллинга появились live-строки
|
||||
// (text/tool/agent/reasoning) — LLM реально работает, сбрасываем idle.
|
||||
if cur := liveSeq.Load(); cur != prevLive {
|
||||
prevLive = cur
|
||||
lastProgress = 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
|
||||
r.logf("opencode(%s) idle %.0fs — abort", agent, r.IdleTimeout.Seconds())
|
||||
if err := c.Abort(ctx, sid); err != nil {
|
||||
r.logf("opencode(%s) abort idle: %v", agent, err)
|
||||
}
|
||||
return &Result{RC: -1, Stdout: "", SessionID: sid}, nil
|
||||
}
|
||||
// hard — общий бюджет от старта запуска.
|
||||
if now.Sub(launch) > r.HardTimeout {
|
||||
r.logf("opencode(%s) hard timeout %.0fs — kill", agent, r.HardTimeout.Seconds())
|
||||
killGroup(proc)
|
||||
killed = true
|
||||
break pollLoop
|
||||
r.logf("opencode(%s) hard timeout %.0fs — abort", agent, r.HardTimeout.Seconds())
|
||||
if err := c.Abort(ctx, sid); err != nil {
|
||||
r.logf("opencode(%s) abort hard: %v", agent, err)
|
||||
}
|
||||
return &Result{RC: -1, Stdout: "", SessionID: sid}, nil
|
||||
}
|
||||
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()
|
||||
|
||||
r.logf("opencode(%s) lines=%d bytes=%d", agent, len(buf), len(out))
|
||||
|
||||
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
|
||||
select {
|
||||
case <-time.After(r.PollInterval):
|
||||
case <-ctx.Done():
|
||||
}
|
||||
}
|
||||
r.logf("opencode(%s) rc=%d", agent, rc)
|
||||
return &Result{RC: rc, Stdout: out, SessionID: sid}, nil
|
||||
}
|
||||
|
||||
// textCount считает число text-частей в assistant-сообщениях (для progress).
|
||||
func (c *Client) textCount(ctx context.Context, sessionID string) (int, error) {
|
||||
msgs, err := c.messages(ctx, sessionID)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
n := 0
|
||||
for _, m := range msgs {
|
||||
if m.Type != "assistant" {
|
||||
continue
|
||||
}
|
||||
for _, p := range m.Content {
|
||||
if p.Type == "text" && p.Text != "" {
|
||||
n++
|
||||
}
|
||||
}
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
|
||||
// messages возвращает сырые сообщения сессии (для поллинга прогресса).
|
||||
func (c *Client) messages(ctx context.Context, sessionID string) ([]sessionMessage, error) {
|
||||
raw, err := c.do(ctx, "GET", "/api/session/"+sessionID+"/message?order=asc&limit=200", "messages", nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var out struct {
|
||||
Data []sessionMessage `json:"data"`
|
||||
}
|
||||
if err := json.Unmarshal(raw, &out); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return out.Data, nil
|
||||
}
|
||||
|
||||
// ResumeDev — запуск dev-агента с resume-fallback. Если resume (sessionID)
|
||||
// падает с rc!=0 (напр. сессия потеряна) — повторяем ОДИН раз свежей сессией
|
||||
// в том же worktree. rc=-1 (kill по таймауту) НЕ триггерит fallback.
|
||||
// падает с rc!=0 — повторяем ОДИН раз свежей сессией в том же каталоге.
|
||||
// rc=-1 (обрыв по таймауту) НЕ триггерит 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)
|
||||
r.logf("dev resume rc=%d — запускаю заново свежей сессией (каталог сохраняю)", 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()
|
||||
}
|
||||
const resumeFallbackNote = "\n\n(Возобновление сессии не удалось; продолжи с учётом уже сделанных изменений в worktree.)"
|
||||
Reference in New Issue
Block a user