diff --git a/cmd/ratatoskr/main.go b/cmd/ratatoskr/main.go index 26a3e18..4a9085d 100644 --- a/cmd/ratatoskr/main.go +++ b/cmd/ratatoskr/main.go @@ -1,11 +1,11 @@ -// Command ratatoskr — оркестратор конвейера Ratatoskr (единственный бинарь). +// Command ratatoskr — оркестратор конвейера Ratatoskr (единый бинарь). // -// Стартовый каркас: точка входа + заглушки пакетов. Полная логика -// (диалог, аналитик, research, feedback, webhook) наполняется по мере -// переезда с Python. +// Запуск: ratatoskr [-config config.yaml] +// config.yaml ищется в CWD, по умолчанию — env/дефолты. package main import ( + "context" "flag" "fmt" "log" @@ -15,7 +15,7 @@ import ( ) func main() { - cfg := flag.String("config", "", "путь к config файлу (по умолчанию — env)") + cfg := flag.String("config", "", "путь к config.yaml (по умолчанию — CWD/config.yaml)") version := flag.Bool("version", false, "показать версию и выйти") flag.Parse() @@ -24,21 +24,20 @@ func main() { return } - // Composition root: читает конфиг (env/config), настраивает DI. a, err := app.New(*cfg) if err != nil { log.Fatalf("app init: %v", err) } + defer func() { + if a.Store != nil { + a.Store.Close() + } + }() - // Пока: заглушка. Здесь будет: - // - запуск Telegram-gateway (диалог) - // - http-сервер webhook (+ /feedback, /status) - // - воркер: обработка issues → opencode-субагенты → PR - _ = a - os.Exit(run()) -} - -func run() int { - fmt.Fprintln(os.Stderr, "ratatoskr-go: каркас готов, конвейер не подключён") - return 0 + log.Print("ratatoskr: запущен") + if err := a.Run(context.Background()); err != nil { + log.Fatalf("app run: %v", err) + } + log.Print("ratatoskr: остановлен") + os.Exit(0) } \ No newline at end of file diff --git a/internal/app/app.go b/internal/app/app.go index 4734758..e09800a 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -1,26 +1,224 @@ -// Package app — composition root: конфигурация, DI, жизненный цикл процесса. +// Package app — composition root, DI, жизненный цикл бинаря. package app -// Config — конфигурация конвейера (env или config-файл). -type Config struct { - GiteaURL string - GiteaToken string - IssueRepo string - // TODO: Telegram, LiteLLM, SQLite (opencode.db), gateway, webhook secret. -} +import ( + "context" + "errors" + "fmt" + "log" + "os" + "os/signal" + "strings" + "syscall" + "time" -// App — собранный конвейер (все зависимости). + "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 - // TODO: *dialog.Controller, *gitea.Client, *opencode.Runner, ... + Config *config.Config + Store *storage.Storage + Router *chat.Router + CoreCtx *core.Core + Worker *worker.Worker + tg *telegram.Channel // сохранена для Run } -// New читает конфиг и собирает приложение. Стартовая заглушка. +// New читает конфиг и собирает все зависимости. +// Не запускает подсистемы (Run). func New(configPath string) (*App, error) { - var c Config - _ = configPath // TODO: загрузка из env/config - // TODO: слой failures/classes ошибок по конвенции проекта. - // Использовать в рантайме (например, лимиты диалога): - c.IssueRepo = "kamelion/ratatoskr-go" // placeholder - return &App{Config: c}, nil + 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() + + 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, + 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, + } + 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) + } } \ No newline at end of file diff --git a/internal/app/app_test.go b/internal/app/app_test.go index 459e4c0..afe3d72 100644 --- a/internal/app/app_test.go +++ b/internal/app/app_test.go @@ -1,15 +1,117 @@ package app -import "testing" +import ( + "context" + "os" + "path/filepath" + "strings" + "testing" +) -// Контракт: process_turn (после внедрения) должен возвращать 3 значения. -// Здесь — placeholder-проверка, что composition root собирается без сбоев. +// TestNew проверяет, что composition root собирается без ошибок. +// Использует :memory: (? или file-based с t.TempDir). func TestNew(t *testing.T) { + tmp := t.TempDir() + dbPath := filepath.Join(tmp, "test.db") + + // Настраиваем минимальные env для обязательных полей + t.Setenv("TG_TOKEN", "test:token") + t.Setenv("TG_CHAT_ID", "12345") + t.Setenv("GITEA_URL", "http://gitea.test") + t.Setenv("GITEA_TOKEN", "test-token") + t.Setenv("RATATOSKR_DB", dbPath) + + // Загружаем без config-файла (дефолты + env) a, err := New("") if err != nil { t.Fatalf("New() err = %v", err) } - if a.Config.IssueRepo == "" { - t.Fatal("IssueRepo пуст — каркас не инициализирован") + if a.Store == nil { + t.Fatal("Store не создан") } + if a.Router == nil { + t.Fatal("Router не создан") + } + if a.CoreCtx == nil { + t.Fatal("Core не создан") + } + if a.Worker == nil { + t.Fatal("Worker не создан") + } + if a.tg == nil { + t.Fatal("Telegram-канал не создан") + } + t.Logf("app OK: db=%s", a.Config.Paths.DB) +} + +// TestNew_MissingToken проверяет ошибку при пустом токене. +func TestNew_MissingToken(t *testing.T) { + tmp := t.TempDir() + dbPath := filepath.Join(tmp, "test.db") + t.Setenv("TG_TOKEN", "") + t.Setenv("TG_CHAT_ID", "12345") + t.Setenv("GITEA_URL", "http://gitea.test") + t.Setenv("GITEA_TOKEN", "test-token") + t.Setenv("RATATOSKR_DB", dbPath) + + _, err := New("") + if err == nil { + t.Fatal("expected error for missing token") + } + t.Logf("got expected error: %v", err) +} + +// TestNew_BadDB проверяет ошибку при пути к БД, который нельзя создать. +func TestNew_BadDB(t *testing.T) { + tmp := t.TempDir() + dbPath := filepath.Join(tmp, "nonexistent", "test.db") // родителя нет + t.Setenv("TG_TOKEN", "test:token") + t.Setenv("TG_CHAT_ID", "12345") + t.Setenv("GITEA_URL", "http://gitea.test") + t.Setenv("GITEA_TOKEN", "test-token") + t.Setenv("RATATOSKR_DB", dbPath) + + _, err := New("") + if err == nil { + t.Fatal("expected error for invalid db path") + } + if !strings.Contains(err.Error(), "db open error") { + t.Errorf("expected 'db open error' in message, got %v", err) + } +} + +// TestNew_CtxCancel проверяет, что Run завершается на отмене контекста. +func TestNew_RunCtxCancel(t *testing.T) { + tmp := t.TempDir() + dbPath := filepath.Join(tmp, "test-runcancel.db") + t.Setenv("TG_TOKEN", "test:token") + t.Setenv("TG_CHAT_ID", "12345") + t.Setenv("GITEA_URL", "http://gitea.test") + t.Setenv("GITEA_TOKEN", "test-token") + t.Setenv("RATATOSKR_DB", dbPath) + + a, err := New("") + if err != nil { + t.Fatalf("New() err = %v", err) + } + + ctx, cancel := context.WithCancel(context.Background()) + cancel() // сразу отменяем + + err = a.Run(ctx) + if err != nil { + t.Fatalf("Run() expected nil, got %v", err) + } + t.Log("Run OK on cancelled context") +} + +// TestMain — дёргает fatalf, нормально что os.Exit не тестируется. +// Просто тест, что main.go компилируется. +func TestMain(t *testing.T) { + // main() не вызываем — он бы os.Exit + _, err := os.Stat("../../cmd/ratatoskr/main.go") + if err != nil { + t.Fatalf("main.go not found: %v", err) + } + t.Log("main.go exists") } \ No newline at end of file diff --git a/internal/app/errors.go b/internal/app/errors.go new file mode 100644 index 0000000..a8f8684 --- /dev/null +++ b/internal/app/errors.go @@ -0,0 +1,17 @@ +package app + +import "errors" + +var ( + // A1 — не загрузился/невалиден конфиг + ErrConfig = errors.New("A1: config error") + + // A2 — не открылась storage-БД + ErrDBOpen = errors.New("A2: db open error") + + // A3 — Telegram-канал упал со смертельной ошибкой (Run вернул ошибку) + ErrChannelFatal = errors.New("A3: channel fatal error") + + // A4 — проблемы при graceful shutdown + ErrShutdown = errors.New("A4: shutdown error") +) \ No newline at end of file diff --git a/internal/config/load.go b/internal/config/load.go index 3e7371b..00136f7 100644 --- a/internal/config/load.go +++ b/internal/config/load.go @@ -13,21 +13,29 @@ import ( // env-переменные (тег env) и валидирует. Если path пустой — ищет config.yaml // в CWD. Ошибка C4 если файл не существует (но не fatal при path=""). func Load(path string) (*Config, error) { + noFile := false if path == "" { path = "config.yaml" } raw, err := os.ReadFile(path) if err != nil { if strings.Contains(err.Error(), "no such file") || strings.Contains(err.Error(), "cannot find") { - return nil, fmt.Errorf("%w: %s", ErrFileNotFound, path) + if path == "config.yaml" { + // Файл не указан и не найден — не фатально, соберём дефолты+env + noFile = true + } else { + return nil, fmt.Errorf("%w: %s", ErrFileNotFound, path) + } + } else { + return nil, fmt.Errorf("%w: read %s: %v", ErrInvalidFormat, path, err) } - return nil, fmt.Errorf("%w: read %s: %v", ErrInvalidFormat, path, err) } - // подстановка ${VAR:-default} - expanded := os.Expand(string(raw), envLookup) cfg := &Config{} - if err := yaml.Unmarshal([]byte(expanded), cfg); err != nil { - return nil, fmt.Errorf("%w: parse %s: %v", ErrInvalidFormat, path, err) + if !noFile { + expanded := os.Expand(string(raw), envLookup) + if err := yaml.Unmarshal([]byte(expanded), cfg); err != nil { + return nil, fmt.Errorf("%w: parse %s: %v", ErrInvalidFormat, path, err) + } } // дефолты для zero-значений applyDefaults(cfg) diff --git a/internal/config/types.go b/internal/config/types.go index 3378830..d2691bf 100644 --- a/internal/config/types.go +++ b/internal/config/types.go @@ -67,6 +67,7 @@ type ChatCfg struct { type PathsCfg struct { Worktree string `yaml:"worktree" default:"/opt/data/ratatoskr/worktrees"` + DB string `yaml:"db" default:"/opt/data/ratatoskr/ratatoskr.db" env:"RATATOSKR_DB"` } // Validate проверяет обязательные поля. Возвращает C1 MissingField diff --git a/internal/storage/tasks.go b/internal/storage/tasks.go index db48fb5..6c65e93 100644 --- a/internal/storage/tasks.go +++ b/internal/storage/tasks.go @@ -87,7 +87,27 @@ func (s *Storage) UpdateTask(ctx context.Context, t *Task) error { return nil } -// ListTasks возвращает список задач с фильтрацией. +// GetActiveTaskByChatID возвращает последнюю не-терминальную задачу чата. +// Терминальные статусы: success, cancelled, aborted, closed. +func (s *Storage) GetActiveTaskByChatID(ctx context.Context, chatID string) (*Task, error) { + t := &Task{} + err := s.db.QueryRowContext(ctx, ` + SELECT id, chat_id, title, goal, repo, why, ac, task_tag, status, created_at, updated_at + FROM tasks + WHERE chat_id = ? AND status NOT IN ('success','cancelled','aborted','closed') + ORDER BY updated_at DESC LIMIT 1`, chatID).Scan( + &t.ID, &t.ChatID, &t.Title, &t.Goal, &t.Repo, + &t.Why, &t.AC, &t.TaskTag, &t.Status, &t.CreatedAt, &t.UpdatedAt) + if err == sql.ErrNoRows { + return nil, fmt.Errorf("%w: no active task for chat %s", ErrNotFound, chatID) + } + if err != nil { + return nil, fmt.Errorf("%w: get active task %s: %w", ErrDB, chatID, err) + } + return t, nil +} + +// ListTasks возвращает задачи по фильтру. func (s *Storage) ListTasks(ctx context.Context, filter TaskFilter) ([]*Task, error) { if filter.Limit <= 0 { filter.Limit = 50