Растущий в один text-парт стрим (text-delta) и reasoning больше не выглядят как зависшая нейронка: idle-таймер сбрасывается по росту числа партов и суммарной длины text/reasoning.
227 lines
9.2 KiB
Go
227 lines
9.2 KiB
Go
package opencode
|
||
|
||
import (
|
||
"context"
|
||
"errors"
|
||
"fmt"
|
||
"io"
|
||
"os"
|
||
"strings"
|
||
"time"
|
||
)
|
||
|
||
// Result — результат запуска opencode-субагента через HTTP API. rc=-1 означает
|
||
// «оборван по таймауту/контексту» (idle/hard): вызывающий НЕ должен ронять
|
||
// задачу, а обязан закоммитить/запушить готовую работу и отправить на ревью
|
||
// (класс O2 Timeout — результат, не ошибка).
|
||
type Result struct {
|
||
RC int
|
||
Stdout string
|
||
SessionID string
|
||
}
|
||
|
||
// Runner — запуск opencode-субагентов через v2 HTTP API serve.
|
||
//
|
||
// Runner ходит к opencode serve через Pool→Client (пути /api/*, см. README,
|
||
// минимальная версия opencode). Промпт отправляется неблокирующе (durable
|
||
// admit), вердикт собирается поллингом новых assistant-сообщений; завершение
|
||
// ответа определяется по схеме «сессия больше не в активных дренажах» + финальное
|
||
// assistant-сообщение.
|
||
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}
|
||
|
||
// Модель по умолчанию из конфига opencode — хардпиним её в сессии, чтобы
|
||
// не зависеть от fallback-логики opencode (класс O5 WARN: если модель не
|
||
// считывается/не задана — предупреждаем и работаем без явного указания).
|
||
model, mErr := ReadModelRef(srv.Config, srv.ConfigDir)
|
||
if mErr != nil {
|
||
r.logf("WARN opencode: не удалось прочитать model из конфига: %v", mErr)
|
||
} else if model == nil {
|
||
r.logf("WARN opencode: в конфиге opencode не задан top-level model — модель не хардпинится (риск fallback)")
|
||
} else {
|
||
r.logf("opencode(%s) model=%s", agent, model)
|
||
}
|
||
|
||
// Сессия: заданная (resume) или новая.
|
||
sid := sessionID
|
||
if sid == "" {
|
||
sid, err = c.CreateSession(ctx, model)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("opencode: create session: %w", err)
|
||
}
|
||
r.logf("opencode(%s) session=%s на %s", agent, sid, srv.Addr())
|
||
}
|
||
|
||
return r.awaitVerdict(ctx, c, model, sid, agent, prompt)
|
||
}
|
||
|
||
// settlePolls — сколько подряд опросов должно подтвердить завершение ответа,
|
||
// прежде чем считать вердикт финальным (устойчивость к гонке между удалением
|
||
// сессии из активных дренажей и финализацией последнего сообщения).
|
||
const settlePolls = 2
|
||
|
||
// awaitVerdict отправляет промпт (неблокирующе) и поллит новые assistant-сообщения,
|
||
// контролируя idle/hard таймауты. Завершение: сессия ушла из активных дренажей
|
||
// И есть новое завершённое assistant-сообщение, стабильное в течение settlePolls
|
||
// опросов. Возвращает вердикт (текст text-партов), либо rc=-1 при таймауте.
|
||
func (r *Runner) awaitVerdict(ctx context.Context, c *Client, model *ModelRef, sid, agent, prompt string) (*Result, error) {
|
||
// admit промпта; граница «новых» сообщений — время создания user-сообщения.
|
||
admittedAt := time.Now().UnixMilli()
|
||
adm, err := c.Prompt(ctx, sid, prompt)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if adm != nil && adm.TimeCreated > 0 {
|
||
admittedAt = adm.TimeCreated
|
||
}
|
||
|
||
// Прогресс = число контент-партов + суммарная длина их текста в новых
|
||
// assistant-сообщениях (progressOf). Рост сбрасывает idle-таймер: LLM
|
||
// стримит (даже в один растущий text-парт) или думает (reasoning) = жив.
|
||
lastParts, lastTextLen := -1, -1
|
||
lastProgress := time.Now()
|
||
launch := time.Now()
|
||
|
||
doneSeen, emptySeen := 0, 0
|
||
|
||
abortAnd := func(rc int, why string) (*Result, error) {
|
||
if err := c.Interrupt(ctx, sid); err != nil {
|
||
r.logf("opencode(%s) interrupt %s: %v", agent, why, err)
|
||
}
|
||
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")
|
||
}
|
||
|
||
msgs, err := c.Messages(ctx, sid)
|
||
if err != nil {
|
||
if ctx.Err() != nil {
|
||
return abortAnd(-1, "ctx")
|
||
}
|
||
var ce *ClientErr
|
||
if errors.As(err, &ce) && ce.Op == "connect" {
|
||
return nil, fmt.Errorf("opencode: %w", err)
|
||
}
|
||
return nil, err
|
||
}
|
||
active, err := c.Active(ctx, sid)
|
||
if err != nil {
|
||
if ctx.Err() != nil {
|
||
return abortAnd(-1, "ctx")
|
||
}
|
||
var ce *ClientErr
|
||
if errors.As(err, &ce) && ce.Op == "connect" {
|
||
return nil, fmt.Errorf("opencode: %w", err)
|
||
}
|
||
return nil, err
|
||
}
|
||
|
||
cur, _ := newestAssistant(msgs, admittedAt)
|
||
parts, textLen := progressOf(msgs, admittedAt)
|
||
if parts != lastParts || textLen != lastTextLen {
|
||
lastProgress = time.Now()
|
||
lastParts, lastTextLen = parts, textLen
|
||
}
|
||
now := time.Now()
|
||
if now.Sub(lastProgress) > r.IdleTimeout {
|
||
r.logf("opencode(%s) idle %.0fs — abort", agent, r.IdleTimeout.Seconds())
|
||
return abortAnd(-1, "idle")
|
||
}
|
||
if now.Sub(launch) > r.HardTimeout {
|
||
r.logf("opencode(%s) hard timeout %.0fs — abort", agent, r.HardTimeout.Seconds())
|
||
return abortAnd(-1, "hard")
|
||
}
|
||
|
||
switch {
|
||
case !active && cur != nil && cur.finished():
|
||
// ответ закончен — ждём стабильности, затем собираем вердикт
|
||
doneSeen++
|
||
emptySeen = 0
|
||
if doneSeen >= settlePolls {
|
||
return r.verdict(model, cur, msgs, admittedAt, sid)
|
||
}
|
||
case !active && cur == nil:
|
||
// сессия завершилась, но нового assistant-сообщения так и нет
|
||
emptySeen++
|
||
if emptySeen >= settlePolls {
|
||
return nil, &ClientErr{Op: "prompt", Err: errors.New("агент не выдал ответ (сессия пуста)")}
|
||
}
|
||
default:
|
||
doneSeen, emptySeen = 0, 0
|
||
}
|
||
|
||
select {
|
||
case <-time.After(r.PollInterval):
|
||
case <-ctx.Done():
|
||
}
|
||
}
|
||
}
|
||
|
||
// verdict собирает финальный результат из новых assistant-сообщений.
|
||
// Проверяет фактическую модель ответа и логирует warning при расхождении
|
||
// с ожидаемой (устойчивость к «не той» модели — класс O5 WARN).
|
||
func (r *Runner) verdict(model *ModelRef, cur *v2Message, msgs []v2Message, since int64, sid string) (*Result, error) {
|
||
if model != nil && cur.Model != nil && (model.ProviderID != cur.Model.ProviderID || model.ID != cur.Model.ID) {
|
||
r.logf("WARN opencode: сессия %s отвечала моделью %s, а не ожидаемой %s — проверь providers в конфиге (v2-схема: provider.api / request, а не npm/options)", sid, cur.Model, model)
|
||
}
|
||
if cur.Error != nil && cur.Error.Message != "" {
|
||
return nil, &ClientErr{Op: "prompt", Err: errors.New(cur.Error.Message)}
|
||
}
|
||
texts := assistantText(msgs, since)
|
||
if len(texts) == 0 {
|
||
return nil, &ClientErr{Op: "prompt", Err: errors.New("нет text-части в ответе")}
|
||
}
|
||
vd := stripFence(strings.Join(texts, "\n"))
|
||
r.logf("opencode вердикт готов (%d байт)", len(vd))
|
||
return &Result{RC: 0, Stdout: vd, SessionID: sid}, nil
|
||
}
|