Files
ratatoskr-go/internal/opencode/runner.go
ki.sagidullin 1459670ce9
Some checks failed
CI / test (pull_request) Failing after 34s
CI / build-and-package (amd64, linux) (pull_request) Successful in 34s
CI / build-and-package (amd64, windows) (pull_request) Successful in 36s
refactor: чистка мёртвого кода, лимит ходов D3, HTML-экранирование и UTF-8 обрезка в Telegram
2026-08-18 23:13:04 +05:00

167 lines
5.8 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package opencode
import (
"context"
"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
Debug bool // отладочные логи API-вызовов (из log.level=debug)
// Заменяемый для тестов:
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, Debug: r.Debug}
// Сессия: заданная (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())
}
// Отправляем промпт (блокирующий Send в горутине; вердикт придёт из него),
// параллельно поллим прогресс и контролируем idle/hard таймауты.
return r.awaitVerdict(ctx, c, sid, agent, prompt)
}
// awaitVerdict запускает блокирующий Send и параллельно поллит прогресс
// (рост числа text-частей = агент жив, сбрасывает idle). Возвращается вердикт
// из ответа Send, либо rc=-1 при idle/hard таймауте (тогда Abort + отмена ctx).
func (r *Runner) awaitVerdict(ctx context.Context, c *Client, sid, agent, prompt string) (*Result, error) {
sendCtx, cancel := context.WithCancel(ctx)
defer cancel()
type sendOut struct {
vd string
err error
}
sendCh := make(chan sendOut, 1)
go func() {
vd, err := c.Send(sendCtx, sid, prompt)
sendCh <- sendOut{vd: vd, err: err}
}()
// Прогресс = сумма text-частей во всех assistant-сообщениях сессии. Рост
// сбрасывает idle-таймер (LLM стримит = жив).
var mu sync.Mutex
lastCount := -1
lastProgress := time.Now()
launch := time.Now()
abortAnd := func(rc int, why string) (*Result, error) {
if err := c.Abort(ctx, sid); err != nil {
r.logf("opencode(%s) abort %s: %v", agent, why, err)
}
cancel()
return &Result{RC: rc, Stdout: "", SessionID: sid}, nil
}
for {
if ctx.Err() != nil {
r.logf("opencode(%s) ctx cancelled — обрыв (rc=-1)", agent)
return abortAnd(-1, "ctx")
}
count, _ := c.textCount(ctx, sid)
mu.Lock()
if count != lastCount {
lastProgress = time.Now()
lastCount = count
}
mu.Unlock()
now := time.Now()
if now.Sub(lastProgress) > r.IdleTimeout {
r.logf("opencode(%s) idle %.0fs — abort", agent, r.IdleTimeout.Seconds())
return abortAnd(-1, "idle")
}
// hard — общий бюджет от старта запуска.
if now.Sub(launch) > r.HardTimeout {
r.logf("opencode(%s) hard timeout %.0fs — abort", agent, r.HardTimeout.Seconds())
return abortAnd(-1, "hard")
}
select {
case out := <-sendCh:
// Send завершился. Ошибка — connect (сервер недоступен) и ctx жив →
// фатально, не таймаут. Если ctx уже отменён — это обрыв, а не ошибка.
if out.err != nil {
var ce *ClientErr
if errors.As(out.err, &ce) && ce.Op == "connect" && ctx.Err() == nil {
return nil, fmt.Errorf("opencode: %w", out.err)
}
if ctx.Err() != nil {
return abortAnd(-1, "ctx")
}
return nil, out.err
}
r.logf("opencode(%s) вердикт готов (%d байт)", agent, len(out.vd))
return &Result{RC: 0, Stdout: out.vd, SessionID: sid}, nil
case <-time.After(r.PollInterval):
case <-ctx.Done():
}
}
}