// 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-), // который вшивается ldflag'ом и используется автообновлением. Здесь номер // поднимается вручную перед каждым релизом/публикацией новой сборки. const Version = "0.2.1" // 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(), Debug: cfg.Log.Debug(), 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 }