Files
ratatoskr-go/internal/app/app.go
Hermes bd3d825738
All checks were successful
CI / test (push) Successful in 46s
CI / build-and-package (amd64, darwin) (push) Successful in 36s
CI / build-and-package (amd64, linux) (push) Successful in 37s
CI / build-and-package (amd64, windows) (push) Successful in 44s
CI / build-and-package (arm64, darwin) (push) Successful in 35s
CI / build-and-package (arm64, linux) (push) Successful in 36s
feat: множественные репозитории (Repos) и клонирование в воркере
- Task.Repos []string (XML-колонка repos, обратная совместимость с repo)
- config: блок git {base_url, token}
- аналитик: ответ repos[], шаблон показывает список
- core: propose без repos → возврат в сбор (E1)
- worker вариант A: один dev из общего cwd, prepareRepos клонирует
  недостающие репо (git clone), validateRepoName (E3), ErrRepoNotGit (E4)
- ошибки E1-E4 в worker/errors.go
2026-08-16 09:18:15 +05:00

255 lines
7.9 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/worker"
)
// App — собранный конвейер.
type App struct {
Config *config.Config
Store *storage.Storage
Router *chat.Router
CoreCtx *core.Core
Worker *worker.Worker
tg *telegram.Channel // сохранена для Run
}
// New читает конфиг и собирает все зависимости.
// Не запускает подсистемы (Run).
func New(configPath 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
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.Print("app: worker started")
// Ожидание сигнала или фатальной ошибки 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)
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)
}
}
// 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
}