app: composition root (A1-A4), Run(ctx), handleIncoming (вариант A), фикс config.Load(""), GetActiveTaskByChatID, main.go обновлён
All checks were successful
build-test / build (push) Successful in 35s

This commit is contained in:
Hermes
2026-08-15 09:51:03 +05:00
parent 4641090398
commit c72b41b2c5
7 changed files with 392 additions and 47 deletions

View File

@@ -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)
}
}