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 сессии
211 lines
7.7 KiB
Go
211 lines
7.7 KiB
Go
package opencode
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"io"
|
||
"os"
|
||
"sync"
|
||
"time"
|
||
)
|
||
|
||
// Result — результат запуска opencode-субагента через HTTP API. rc=-1 означает
|
||
// «оборван по таймауту/контексту» (idle/hard): вызывающий НЕ должен ронять
|
||
// задачу, а обязан закоммитить/запушить готовую работу и отправить на ревью
|
||
// (класс O2 Timeout — результат, не ошибка).
|
||
type Result struct {
|
||
RC int
|
||
Stdout string
|
||
SessionID string
|
||
}
|
||
|
||
// Runner — запуск opencode-субагентов через HTTP API serve.
|
||
//
|
||
// Полный переход на API: Runner ходит к opencode serve через Pool→Client
|
||
// (нет spawn-модели, нет NDJSON). Агент идёт в сервер пула для своего каталога
|
||
// (в нём запущен serve → он его project).
|
||
type Runner struct {
|
||
Pool *Pool // пул serve-серверов (обязательный)
|
||
IdleTimeout time.Duration
|
||
HardTimeout time.Duration
|
||
PollInterval time.Duration
|
||
|
||
// Заменяемый для тестов:
|
||
Stdout io.Writer // диагностика (лог), по умолчанию os.Stderr
|
||
}
|
||
|
||
func (r *Runner) defaults() {
|
||
if r.IdleTimeout == 0 {
|
||
r.IdleTimeout = 5 * 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...)
|
||
}
|
||
|
||
// 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()
|
||
if r.Pool == nil {
|
||
return nil, fmt.Errorf("opencode: Pool не задан (API-режим обязателен)")
|
||
}
|
||
|
||
srv, err := r.Pool.Ensure(ctx, cwd)
|
||
if err != nil {
|
||
return nil, 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())
|
||
}
|
||
|
||
// Отправляем промпт (неблокирующий — сервер начинает выполнение).
|
||
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
|
||
lastCount := -1
|
||
lastProgress := time.Now()
|
||
launch := time.Now()
|
||
|
||
for {
|
||
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
|
||
}
|
||
|
||
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()
|
||
if now.Sub(lastProgress) > r.IdleTimeout {
|
||
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 — 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
|
||
}
|
||
|
||
select {
|
||
case <-time.After(r.PollInterval):
|
||
case <-ctx.Done():
|
||
}
|
||
}
|
||
}
|
||
|
||
// 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 — повторяем ОДИН раз свежей сессией в том же каталоге.
|
||
// 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 {
|
||
return res, false
|
||
}
|
||
if res.RC != 0 && res.RC != -1 && sessionID != "" {
|
||
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.)" |