// 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) } // Относительные пути (db, worktree) — рядом с .exe, а не от CWD запуска. cfg.ResolveExePaths() // 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 }