Files
ratatoskr-go/internal/opencode/pool.go
Hermes 733e63339a
Some checks failed
CI / test (push) Successful in 40s
CI / build-and-package (amd64, linux) (push) Successful in 35s
CI / build-and-package (amd64, windows) (push) Failing after 26s
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 сессии
2026-08-18 13:36:54 +05:00

234 lines
6.6 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"
"fmt"
"log"
"net"
"path/filepath"
"sync"
)
// Pool — контроль над пулом opencode serve-процессов (по одному на каталог).
//
// Ленивый: сервер для каталога поднимается при первом запросе (Ensure), кроме
// служебного root-сервера (EnsureRoot), который живёт с момента старта app.
// Завершение задачи снимает поднятые серверы кроме root'a (ReleaseTask).
//
// Каждый Server слушает свой порт (basePort + сдвиг), запускается в своей
// директории → каждая сессия API привязана к правильному project-каталогу.
type Pool struct {
Bin string
Config string
ConfigDir string
DBPath string
Host string
BasePort int
Password string
rootDir string // каталог служебного сервера
root *Server
ctx context.Context // базовый ctx для всех serve; живёт, пока пул активен
cancel context.CancelFunc
mu sync.Mutex
segs map[string]*Server // dir → сервер (root тоже здесь)
used map[int]bool // занятые порты
next int // следующий кандидат порта
}
// NewPool создаёт пул. rootDir помечен как служебный (не снимается ReleaseTask).
func NewPool(rootDir string) *Pool {
return &Pool{
Host: "127.0.0.1",
BasePort: 4096,
rootDir: rootDir,
segs: make(map[string]*Server),
used: make(map[int]bool),
next: 4096,
}
}
// startMonitored поднимает сервер и запускает его Run-перезапуск (reaper).
// Наследует базовый ctx пула: Serve живёт, пока жив пул.
func (p *Pool) startMonitored(ctx context.Context, s *Server) error {
if p.ctx == nil {
sctx, cancel := context.WithCancel(ctx)
p.ctx, p.cancel = sctx, cancel
}
if err := s.Start(p.ctx); err != nil {
return err
}
go s.Run(p.ctx)
return nil
}
// EnsureRoot поднимает служебный сервер в rootDir (идемпотентен).
func (p *Pool) EnsureRoot(ctx context.Context) error {
p.mu.Lock()
defer p.mu.Unlock()
if p.root != nil {
return nil
}
s := &Server{
Bin: p.Bin,
Config: p.Config,
ConfigDir: p.ConfigDir,
DBPath: p.DBPath,
Host: p.Host,
Password: p.Password,
Dir: p.rootDir,
}
if err := p.assign(s); err != nil {
return err
}
if err := p.startMonitored(ctx, s); err != nil {
return fmt.Errorf("opencode serve (root): %w", err)
}
p.root = s
p.segs[p.rootDir] = s
log.Printf("opencode: root serve up at %s (dir %s)", s.Addr(), p.rootDir)
return nil
}
// Ensure гарантирует наличие сервера для каталога dir (лениво).
// Возвращает сервер; root-сервер для rootDir возвращается как есть.
func (p *Pool) Ensure(ctx context.Context, dir string) (*Server, error) {
p.mu.Lock()
if s, ok := p.segs[dir]; ok {
p.mu.Unlock()
return s, nil
}
abs := filepath.Clean(dir)
s := &Server{
Bin: p.Bin,
Config: p.Config,
ConfigDir: p.ConfigDir,
DBPath: p.DBPath,
Host: p.Host,
Password: p.Password,
Dir: abs,
}
if err := p.assign(s); err != nil {
p.mu.Unlock()
return nil, err
}
p.segs[abs] = s
p.mu.Unlock()
if err := p.startMonitored(ctx, s); err != nil {
p.mu.Lock()
delete(p.segs, abs)
p.releasePort(s.Port)
p.mu.Unlock()
return nil, fmt.Errorf("opencode serve (%s): %w", abs, err)
}
log.Printf("opencode: serve up at %s (dir %s)", s.Addr(), abs)
return s, nil
}
// RegisterExternal регистрирует внешний (уже запущенный) сервер для каталога
// dir. Полезно, когда serve поднят вне пула (в т.ч. в тестах): Ensure вернёт
// его без spawn'а. url — полный адрес, по которому Runner ходит через API.
func (p *Pool) RegisterExternal(dir, url string) {
p.mu.Lock()
defer p.mu.Unlock()
abs := filepath.Clean(dir)
p.segs[abs] = &Server{URL: url, Host: p.Host, PollInterval: 0}
if p.rootDir != "" && abs == p.rootDir {
p.root = p.segs[abs]
}
}
// ServerFor возвращает сервер для каталога (без поднятия). ok=false если нет.
func (p *Pool) ServerFor(dir string) (*Server, bool) {
p.mu.Lock()
defer p.mu.Unlock()
s, ok := p.segs[filepath.Clean(dir)]
return s, ok
}
// ReleaseTask закрывает все серверы пула, кроме служебного root. Вызывается
// при завершении задачи.
func (p *Pool) ReleaseTask() {
p.mu.Lock()
var toClose []*Server
for dir, s := range p.segs {
if p.root != nil && dir == p.rootDir {
continue // служебный не снимаем
}
toClose = append(toClose, s)
delete(p.segs, dir)
p.releasePort(s.Port)
}
p.mu.Unlock()
for _, s := range toClose {
log.Printf("opencode: closing serve %s (release task)", s.Addr())
s.Close()
}
}
// Close закрывает все серверы пула, включая root. Идемпотентен.
func (p *Pool) Close() {
p.mu.Lock()
toClose := make([]*Server, 0, len(p.segs))
for dir, s := range p.segs {
toClose = append(toClose, s)
delete(p.segs, dir)
p.releasePort(s.Port)
}
p.root = nil
if p.cancel != nil {
p.cancel()
p.cancel = nil
}
p.mu.Unlock()
for _, s := range toClose {
s.Close()
}
}
// assign выделяет свободный порт и проставляет его серверу.
func (p *Pool) assign(s *Server) error {
port, err := p.allocPort()
if err != nil {
return err
}
s.Port = port
return nil
}
// allocPort находит свободный порт начиная с next, коммитит его.
func (p *Pool) allocPort() (int, error) {
for i := 0; i < 100; i++ {
port := p.next
p.next++
if p.used[port] {
continue
}
if !portFree(p.Host, port) {
p.used[port] = true
continue
}
p.used[port] = true
return port, nil
}
return 0, fmt.Errorf("opencode: нет свободных портов в диапазоне")
}
func (p *Pool) releasePort(port int) {
if port != 0 {
delete(p.used, port)
}
}
// portFree проверяет, свободен ли порт (bind probe).
func portFree(host string, port int) bool {
l, err := net.Listen("tcp", fmt.Sprintf("%s:%d", host, port))
if err != nil {
return false
}
l.Close()
return true
}