Files
ratatoskr-go/internal/app/app.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

559 lines
22 KiB
Go
Raw 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 app — composition root, DI, жизненный цикл бинаря.
package app
import (
"context"
"errors"
"fmt"
"log"
"os"
"os/signal"
"path/filepath"
"strings"
"syscall"
"time"
"github.com/kamelion/ratatoskr-go/internal/agents"
"github.com/kamelion/ratatoskr-go/internal/analyst"
"github.com/kamelion/ratatoskr-go/internal/chat"
"github.com/kamelion/ratatoskr-go/internal/chat/telegram"
"github.com/kamelion/ratatoskr-go/internal/config"
"github.com/kamelion/ratatoskr-go/internal/core"
"github.com/kamelion/ratatoskr-go/internal/opencode"
"github.com/kamelion/ratatoskr-go/internal/storage"
"github.com/kamelion/ratatoskr-go/internal/update"
"github.com/kamelion/ratatoskr-go/internal/worker"
)
// packageOwner — владелец Gitea-пакета, из которого берутся обновления.
// Жёсткая константа: владелец пакета совпадает с владельцем репозитория и
// создателем токена read:package (kamelion). Не требует конфигурации.
const packageOwner = "kamelion"
// Version — семантический номер версии приложения в формате major.minor.patch
// (например "0.1.0"). Меняется вручную при выпуске новых изменений; НЕ должен
// вшиваться ldflag'ом или генерироваться автоматически.
//
// Правила ручного инкремента (когда и какую часть номера увеличивать):
// - patch (0.1.0 → 0.1.1): исправление багов и мелкие правки, новая
// функциональность не добавляется (обратно-совместимые изменения).
// - minor (0.1.0 → 0.2.0): появляется новая (обратно-совместимая)
// функциональность; patch при этом сбрасывается в 0.
// - major (0.1.0 → 1.0.0): несовместимые изменения API/поведения или крупные
// релизы; minor и patch сбрасываются в 0.
//
// Пока продукт не стабилен, major держим на 0 → версии идут 0.x.y
// (минорные правки с повышением minor, исправления — с повышением patch).
//
// ВАЖНО: это СЕМАНТИЧЕСКАЯ версия приложения (для людей и диагностики), её
// не следует путать с build-идентификатором `main.version` (commit-<sha7>),
// который вшивается ldflag'ом и используется автообновлением. Здесь номер
// поднимается вручную перед каждым релизом/публикацией новой сборки.
const Version = "0.2.0"
// App — собранный конвейер.
type App struct {
Config *config.Config
Store *storage.Storage
Router *chat.Router
CoreCtx *core.Core
Worker *worker.Worker
Updater *update.Updater
tg *telegram.Channel // сохранена для Run
pool *opencode.Pool // пул opencode serve-серверов (API-режим)
}
// New читает конфиг и собирает все зависимости.
// Не запускает подсистемы (Run).
// version — вшитая версия бинаря (ldflag -X main.version).
// updateToken — вшитый токен read:package для авто-обновления
// (ldflag -X main.updateToken); имеет приоритет над update.token из конфига.
func New(configPath, version, updateToken string) (*App, error) {
cfg, err := config.Load(configPath)
if err != nil {
return nil, fmt.Errorf("%w: %v", ErrConfig, err)
}
if err := cfg.Validate(); err != nil {
return nil, fmt.Errorf("%w: %v", ErrConfig, err)
}
// Относительные пути (db, worktree) — рядом с .exe, а не от CWD запуска.
cfg.ResolveExePaths()
// Каталог worktree может отсутствовать на свежей машине. Decider (analyst) и
// воркер запускают opencode с cwd=worktree, а git clone не создаёт родительский
// каталог — поэтому создаём его заранее, до любых запусков субагентов.
if cfg.Paths.Worktree != "" {
if err := os.MkdirAll(cfg.Paths.Worktree, 0o755); err != nil {
return nil, fmt.Errorf("app: создать каталог worktree %s: %w", cfg.Paths.Worktree, err)
}
}
// SQLite для задач и трасс
ctx := context.Background()
// Агенты opencode: если config_dir не задан — используем ./agents рядом с бинарём
if err := ensureAgentsDir(cfg); err != nil {
return nil, fmt.Errorf("%w: %v", ErrConfig, err)
}
store, err := storage.Open(ctx, cfg.Paths.DB)
if err != nil {
return nil, fmt.Errorf("%w: %v", ErrDBOpen, err)
}
log.Printf("app: db opened %s", cfg.Paths.DB)
// OpenCode: пул serve-процессов (по одному на каталог) + API-runner.
// Служебный root-сервер (worktree) живёт всё время app; остальные лениво.
ocPool := opencode.NewPool(cfg.Paths.Worktree)
ocPool.Bin = cfg.OpenCode.Bin
ocPool.Config = cfg.OpenCode.Config
ocPool.ConfigDir = cfg.OpenCode.ConfigDir
ocPool.DBPath = cfg.OpenCode.DBPath
ocPool.Host = cfg.OpenCode.Serve.Hostname
ocPool.BasePort = cfg.OpenCode.Serve.Port
ocPool.Password = cfg.OpenCode.Serve.Password
// OpenCode runner — один на аналитика и воркер
ocRunner := &opencode.Runner{
Pool: ocPool,
IdleTimeout: cfg.OpenCode.IdleTimeout.Duration(),
HardTimeout: cfg.OpenCode.HardTimeout.Duration(),
PollInterval: cfg.OpenCode.PollMs.Duration(),
Stdout: os.Stderr,
}
// Analyst (Decider)
analystCtx := &analyst.Analyst{
Runner: ocRunner,
Worktree: cfg.Paths.Worktree,
Agent: "analyst",
}
// Core — ядро машины состояний
coreCtx := core.New(store, analystCtx)
// Live-журнал живых сессий: воркер пишет шаги агента, /status N их читает.
live := opencode.NewLiveRegistry()
// Router — единый диспетчер входящих из всех каналов
a := &App{
Config: cfg,
Store: store,
CoreCtx: coreCtx,
pool: ocPool,
}
router := chat.NewRouter(a.handleIncoming)
// Telegram-канал
tg := telegram.New(cfg.Telegram.Token, cfg.Chat.PollInterval.Duration())
a.tg = tg
if err := router.Attach(tg); err != nil {
store.Close()
return nil, fmt.Errorf("attach telegram: %w", err)
}
// Worker — polling-планировщик dev-агента
w := &worker.Worker{
Store: store,
Runner: ocRunner,
Worktree: cfg.Paths.Worktree,
Agent: "dev",
Interval: 5 * time.Second,
MaxJobs: 2,
GitBaseURL: cfg.Git.BaseURL,
GitToken: cfg.Git.Token,
Live: live,
Notify: a, // авто-уведомления владельцу задачи через Router
}
a.Router = router
a.Worker = w
// /status N читает живое состояние сессии через liveProbe-адаптер.
coreCtx.Live = &liveProbe{reg: live}
// Auto-обновление: .new/.old рядом с бинарником (Dir пуст → binDir() от os.Executable).
// Токен: вшитый updateToken приоритетнее update.token из конфига.
udToken := updateToken
if udToken == "" {
udToken = cfg.Update.Token
}
ud := &update.Updater{
BaseURL: cfg.Update.BaseURL,
Owner: packageOwner,
Package: cfg.Update.Package,
Token: udToken,
CurrentVersion: version,
}
a.Updater = ud
return a, nil
}
// Run запускает все подсистемы и блокируется до сигнала завершения.
// Возвращает: nil при graceful shutdown, A3 при фатальной ошибке Telegram.
func (a *App) Run(ctx context.Context) error {
ctx, cancel := context.WithCancel(ctx)
defer cancel()
// Уже отменённый контекст — не поднимаем подсистемы, graceful shutdown сразу.
if ctx.Err() != nil {
log.Print("app: context already cancelled, skipped start")
return nil
}
// opencode serve: поднимаем служебный корневой сервер (worktree) до старта
// воркера, остальные каталоги — лениво. При неудаче — не стартуем.
if err := a.pool.EnsureRoot(ctx); err != nil {
return fmt.Errorf("opencode: %w", err)
}
defer a.pool.Close()
// Канал для проверки Telegram-ошибки (горутина оборачивает Run)
tgErr := make(chan error, 1)
// Telegram: long-poll цикл в горутине
go func() {
tg := a.telegramChannel()
log.Print("app: telegram poll started")
tgErr <- tg.Run(ctx)
}()
// Worker: poll-цикл (неблокирующий — стартует свою горутину)
a.Worker.Start(ctx)
log.Printf("app: worker started")
// Авто-проверка обновления (только уведомление владельца; замена — по /update)
a.startAutoCheck(ctx)
// Ожидание сигнала или фатальной ошибки Telegram
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
select {
case <-ctx.Done():
log.Print("app: context cancelled")
return nil
case sig := <-sigCh:
log.Printf("app: signal %s — shutting down", sig)
cancel()
return nil
case err := <-tgErr:
if err != nil && !errors.Is(err, context.Canceled) {
return fmt.Errorf("%w: %v", ErrChannelFatal, err)
}
return nil
}
}
// telegramChannel возвращает Telegram-канал.
func (a *App) telegramChannel() *telegram.Channel { return a.tg }
// handleIncoming — колбэк Router на каждое входящее сообщение.
// Реализует политику «одна активная задача на чат» (вариант A).
func (a *App) handleIncoming(inc chat.Incoming) {
ctx := context.Background()
uid := inc.UserID
text := strings.TrimSpace(inc.Msg.Text)
if text == "" {
return
}
log.Printf("app: incoming from %s: %q", uid, text)
// update/status — команды бинаря вне машины состояний задач.
// /status перехватываем только без числового аргумента (/status N — статус задачи).
// /help — справка по всем командам (тоже команда бинаря, вне задач).
if cmdName(text) == "/help" ||
cmdName(text) == "/update" ||
(cmdName(text) == "/status" && !hasArg(text)) {
a.handleBinaryCommand(ctx, uid, text)
return
}
task, err := a.Store.GetActiveTaskByChatID(ctx, string(uid))
if err != nil && !errors.Is(err, storage.ErrNotFound) {
log.Printf("app: get active task: %v", err)
a.send(ctx, uid, "Ошибка базы данных")
return
}
// Нет активной задачи
if task == nil {
switch {
case text == "/start":
id, err := a.Store.CreateTask(ctx, &storage.Task{ChatID: string(uid)})
if err != nil {
log.Printf("app: create task: %v", err)
a.send(ctx, uid, "Ошибка создания задачи")
return
}
task, err = a.Store.GetTask(ctx, id)
if err != nil {
log.Printf("app: get fresh task: %v", err)
a.send(ctx, uid, "Ошибка создания задачи")
return
}
result, err := a.CoreCtx.ProcessTurn(ctx, task.ID, "/start")
if err != nil {
log.Printf("app: process /start #%d: %v", task.ID, err)
a.send(ctx, uid, "Ошибка обработки /start")
return
}
a.send(ctx, uid, result.Reply)
case strings.HasPrefix(text, "/status"),
strings.HasPrefix(text, "/retry"),
strings.HasPrefix(text, "/continue"):
result, err := a.CoreCtx.ProcessTurn(ctx, 0, text)
if err != nil {
log.Printf("app: process cross-task %q: %v", text, err)
a.send(ctx, uid, "Ошибка обработки")
return
}
a.send(ctx, uid, result.Reply)
default:
a.send(ctx, uid, "Нет активной задачи. /start — создать новую.")
}
return
}
// Активная задача есть, но воркер её выполняет (running). Обычный ввод и команды
// редактирования НЕ должны переводить задачу обратно в collecting — это невалидный
// переход S3 (running → collecting). Разрешено лишь терминальное /cancel и /status N.
if task.Status == storage.StatusRunning {
switch {
case cmdName(text) == "/cancel":
// переходим на терминальный статус
case cmdName(text) == "/status" && hasArg(text):
// запрос статуса конкретной задачи — безопасно, не трогает машину состояний
default:
a.send(ctx, uid, "⏳ Задача ещё выполняется воркером. Дождитесь результата, или /cancel чтобы остановить.")
return
}
}
// Активная задача есть — обработка
result, err := a.CoreCtx.ProcessTurn(ctx, task.ID, text)
if err != nil {
log.Printf("app: process #%d: %v", task.ID, err)
a.send(ctx, uid, "Ошибка обработки")
return
}
a.send(ctx, uid, result.Reply)
}
// send отправляет сообщение пользователю через маршрут.
func (a *App) send(ctx context.Context, uid chat.UserID, text string) {
if err := a.Router.Send(ctx, uid, chat.Message{Text: text}); err != nil {
log.Printf("app: send to %s: %v", uid, err)
}
}
// Notify реализует worker.Notifier: авто-уведомление владельцу задачи через
// chat.Router.Send (переходы статусов и хендоффы dev↔reviewer со стороны воркера).
// Router nil (тесты без Router / ранняя инициализация) — тихо пропускаем.
func (a *App) Notify(ctx context.Context, taskID int64, chatID, text string) error {
if a.Router == nil {
return nil
}
return a.Router.Send(ctx, chat.UserID(chatID), chat.Message{Text: text})
}
// cmdName извлекает команду (первое слово до пробела, нижний регистр).
func cmdName(text string) string {
s := strings.TrimSpace(text)
if i := strings.IndexByte(s, ' '); i > 0 {
s = s[:i]
}
return strings.ToLower(s)
}
// hasArg возвращает true, если после команды есть аргумент (например /status 5).
func hasArg(text string) bool {
s := strings.TrimSpace(text)
i := strings.IndexByte(s, ' ')
return i > 0 && strings.TrimSpace(s[i:]) != ""
}
// ownerUID — куда слать авто-уведомления об обновлении.
func (a *App) ownerUID() chat.UserID {
if a.Config.Telegram.ChatID != "" {
return chat.UserID(a.Config.Telegram.ChatID)
}
return ""
}
// startAutoCheck запускает фоновую проверку наличия обновления.
// Только уведомляет владельца; скачивание/замена — по команде /update.
func (a *App) startAutoCheck(ctx context.Context) {
if !a.Config.Update.Enabled {
return
}
if a.Updater == nil || a.Updater.BaseURL == "" {
log.Print("app: update авто-проверка выключена (не заполнен update.base_url)")
return
}
interval := a.Config.Update.CheckInterval.Duration()
if interval <= 0 {
interval = 24 * time.Hour
}
go func() {
// первая проверка сразу после старта
a.autoCheckOnce(ctx)
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
a.autoCheckOnce(ctx)
}
}
}()
}
// autoCheckOnce — одна проверка + уведомление владельца.
func (a *App) autoCheckOnce(ctx context.Context) {
res := a.Updater.Check(ctx)
if res.Err != nil {
log.Printf("update: auto-check: %v", res.Err)
return
}
if !res.UpdateAvailable {
return
}
uid := a.ownerUID()
if uid == "" {
log.Printf("update: доступна версия %s (нет chat_id для уведомления)", res.Version)
return
}
msg := fmt.Sprintf("⚡ Доступна новая версия ratatoskr **%s** (у меня %s). Выполни /update для применения.",
res.Version, a.Updater.CurrentVersion)
a.send(ctx, uid, msg)
}
// handleBinaryCommand — маршрутизация бинарных команд: /help, /update, /status.
func (a *App) handleBinaryCommand(ctx context.Context, uid chat.UserID, text string) {
switch cmdName(text) {
case "/status":
a.handleStatusCmd(ctx, uid)
case "/update":
a.handleUpdateCmd(ctx, uid)
case "/help":
a.handleHelpCmd(ctx, uid)
}
}
// helpTextFor — полный текст команды /help с номером версии приложения.
func helpTextFor(version string) string {
return "⚡ **Ratatoskr** — помогу продумать задачу и запущу агента.\n\n" +
"**Команды:**\n" +
"`/start` — новая задача\n" +
"`/cancel` — отменить текущую\n" +
"`/skip` — хватит вопросов, предложить черновик\n" +
"`/retry N` — перезапустить задачу `N`\n" +
"`/status N` — статус задачи `N`\n" +
"`/continue N` — продолжить задачу `N`\n" +
"`/status` — статус бинаря и обновления\n" +
"`/update` — применить обновление\n" +
"`/help` — эта справка\n\n" +
"Версия: " + version + "\n\n" +
"Просто опишите задачу — я помогу её продумать."
}
// handleHelpCmd отвечает списком команд и номером версии приложения.
func (a *App) handleHelpCmd(ctx context.Context, uid chat.UserID) {
a.send(ctx, uid, helpTextFor(Version))
}
// handleStatusCmd отвечает текущей версией и результатом последней проверки.
func (a *App) handleStatusCmd(ctx context.Context, uid chat.UserID) {
if a.Config.Update.Enabled && a.Updater != nil {
res := a.Updater.Check(ctx)
if res.Err != nil {
a.send(ctx, uid, fmt.Sprintf("ratatoskr **%s** (проверка обновления: %v)", a.version(), res.Err))
return
}
if res.UpdateAvailable {
a.send(ctx, uid, fmt.Sprintf("ratatoskr **%s** — доступно обновление до **%s** (`/update`)", a.version(), res.Version))
return
}
a.send(ctx, uid, fmt.Sprintf("ratatoskr **%s** — актуальная версия (latest %s)", a.version(), res.Version))
return
}
a.send(ctx, uid, fmt.Sprintf("ratatoskr **%s** (авто-обновление выключено)", a.version()))
}
// version возвращает текущую вшитую версию.
func (a *App) version() string {
if a.Updater != nil && a.Updater.CurrentVersion != "" {
return a.Updater.CurrentVersion
}
return "dev"
}
// handleUpdateCmd — команда /update: Check → Download → Verify → Swap/Restart.
// При активной задаче в чате — требует подтверждения, чтобы не прерывать работу.
func (a *App) handleUpdateCmd(ctx context.Context, uid chat.UserID) {
if a.Updater == nil || a.Updater.BaseURL == "" {
a.send(ctx, uid, "Обновление не настроено (нужен update.base_url в config).")
return
}
if !a.Config.Update.Enabled {
a.send(ctx, uid, "Авто-обновление выключено (update.enabled: false).")
return
}
res := a.Updater.Check(ctx)
if res.Err != nil {
a.send(ctx, uid, "Не удалось проверить обновление: "+res.Err.Error())
return
}
if !res.UpdateAvailable {
a.send(ctx, uid, fmt.Sprintf("Уже актуальная версия (**%s**).", res.Version))
return
}
a.send(ctx, uid, fmt.Sprintf("Скачиваю **%s**…", res.Version))
file, err := a.Updater.Download(ctx, res.Version)
if err != nil {
a.send(ctx, uid, "Ошибка скачивания: "+err.Error())
return
}
if err := a.Updater.Verify(ctx, res.Version, file); err != nil {
a.send(ctx, uid, "Обновление отклонено (контрольная сумма): "+err.Error())
return
}
// попытка мгновенного свапа + рестарта (POSIX); на Windows — свап при старте
if err := a.Updater.SwapAndRestart(file); err != nil {
// SwapAndRestart.Rewind неexec ищет момент: если Windows — говорим о рестарте
a.send(ctx, uid, "Обновление применено при следующем запуске: "+err.Error())
return
}
// сюда не возвращаемся — SwapAndRestart завершил процесс (os.Exit)
}
// ensureAgentsDir определяет каталог с агентами opencode и распаковывает
// туда встроенных агентов (go:embed). Если config_dir не задан — использует
// ./agents рядом с бинарём. Встроенные агенты перезаписываются (всегда актуальны).
//
// idempotent: вызывается только из New.
func ensureAgentsDir(cfg *config.Config) error {
if cfg.OpenCode.ConfigDir == "" {
exe, err := os.Executable()
if err != nil {
return fmt.Errorf("resolve executable: %w", err)
}
// ./agents рядом с бинарём
cfg.OpenCode.ConfigDir = filepath.Join(filepath.Dir(exe), "agents")
}
if err := agents.WriteTo(cfg.OpenCode.ConfigDir); err != nil {
return fmt.Errorf("write agents to %s: %w", cfg.OpenCode.ConfigDir, err)
}
log.Printf("app: agents ensured in %s", cfg.OpenCode.ConfigDir)
return nil
}