All checks were successful
CI / test (push) Successful in 43s
CI / build-and-package (amd64, darwin) (push) Successful in 40s
CI / build-and-package (amd64, linux) (push) Successful in 39s
CI / build-and-package (amd64, windows) (push) Successful in 41s
CI / build-and-package (arm64, darwin) (push) Successful in 42s
CI / build-and-package (arm64, linux) (push) Successful in 41s
225 lines
6.6 KiB
Go
225 lines
6.6 KiB
Go
// Package app — composition root, DI, жизненный цикл бинаря.
|
||
package app
|
||
|
||
import (
|
||
"context"
|
||
"errors"
|
||
"fmt"
|
||
"log"
|
||
"os"
|
||
"os/signal"
|
||
"strings"
|
||
"syscall"
|
||
"time"
|
||
|
||
"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()
|
||
|
||
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,
|
||
}
|
||
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)
|
||
}
|
||
} |