Files
ratatoskr-go/internal/app/app.go
Hermes 91e91ac3ce
All checks were successful
CI / test (push) Successful in 47s
CI / build-and-package (amd64, linux) (push) Successful in 41s
CI / build-and-package (amd64, windows) (push) Successful in 41s
feat: команда /help — справка по всем командам бота
2026-08-16 20:05:11 +05:00

466 lines
16 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"
// 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
}
// 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)
}
// 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 runner — один на аналитика и воркер
ocRunner := &opencode.Runner{
Bin: cfg.OpenCode.Bin,
DBPath: cfg.OpenCode.DBPath,
Config: cfg.OpenCode.Config,
ConfigDir: cfg.OpenCode.ConfigDir,
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)
// Router — единый диспетчер входящих из всех каналов
a := &App{
Config: cfg,
Store: store,
CoreCtx: coreCtx,
}
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,
}
a.Router = router
a.Worker = w
// 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()
// Канал для проверки 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
}
// Активная задача есть — обработка
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)
}
}
// 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)
}
}
// helpText — текст команды /help.
const helpText = "⚡ **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" +
"Просто опишите задачу — я помогу её продумать."
// handleHelpCmd отвечает списком команд.
func (a *App) handleHelpCmd(ctx context.Context, uid chat.UserID) {
a.send(ctx, uid, helpText)
}
// 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
}