Compare commits
6 Commits
bdd9b41e56
...
feat/48f52
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
07bae203c3 | ||
|
|
43266ea04f | ||
|
|
f12ff04966 | ||
| a8a2de12c3 | |||
|
|
1b24468abc | ||
|
|
d41c0d172f |
@@ -52,7 +52,9 @@ func main() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Print("ratatoskr: запущен")
|
// Строка с номером версии при старте: семантическая версия (app.Version)
|
||||||
|
// + build-идентификатор (main.version, вшитый ldflag'ом) для диагностики.
|
||||||
|
log.Printf("ratatoskr: запущен (version %s, build %s)", app.Version, version)
|
||||||
if err := a.Run(context.Background()); err != nil {
|
if err := a.Run(context.Background()); err != nil {
|
||||||
log.Fatalf("app run: %v", err)
|
log.Fatalf("app run: %v", err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -30,6 +30,27 @@ import (
|
|||||||
// создателем токена read:package (kamelion). Не требует конфигурации.
|
// создателем токена read:package (kamelion). Не требует конфигурации.
|
||||||
const packageOwner = "kamelion"
|
const packageOwner = "kamelion"
|
||||||
|
|
||||||
|
// Version — семантический номер версии приложения в формате major.minor.patch
|
||||||
|
// (например "0.1.0"). Меняется вручную при выпуске новых изменений; НЕ должен
|
||||||
|
// вшиваться ldflag'ом или генерироваться автоматически.
|
||||||
|
//
|
||||||
|
// Правила ручного инкремента (когда и какую часть номера увеличивать):
|
||||||
|
// - patch (0.1.0 → 0.1.1): исправление багов и мелкие правки, новая
|
||||||
|
// функциональность не добавляется (обратно-совместимые изменения).
|
||||||
|
// - minor (0.1.0 → 0.2.0): появляется новая (обратно-совместимая)
|
||||||
|
// функциональность; patch при этом сбрасывается в 0.
|
||||||
|
// - major (0.1.0 → 1.0.0): несовместимые изменения API/поведения или крупные
|
||||||
|
// релизы; minor и patch сбрасываются в 0.
|
||||||
|
//
|
||||||
|
// Пока продукт не стабилен, major держим на 0 → версии идут 0.x.y
|
||||||
|
// (минорные правки с повышением minor, исправления — с повышением patch).
|
||||||
|
//
|
||||||
|
// ВАЖНО: это СЕМАНТИЧЕСКАЯ версия приложения (для людей и диагностики), её
|
||||||
|
// не следует путать с build-идентификатором `main.version` (commit-<sha7>),
|
||||||
|
// который вшивается ldflag'ом и используется автообновлением. Здесь номер
|
||||||
|
// поднимается вручную перед каждым релизом/публикацией новой сборки.
|
||||||
|
const Version = "0.1.0"
|
||||||
|
|
||||||
// App — собранный конвейер.
|
// App — собранный конвейер.
|
||||||
type App struct {
|
type App struct {
|
||||||
Config *config.Config
|
Config *config.Config
|
||||||
@@ -133,6 +154,7 @@ func New(configPath, version, updateToken string) (*App, error) {
|
|||||||
GitBaseURL: cfg.Git.BaseURL,
|
GitBaseURL: cfg.Git.BaseURL,
|
||||||
GitToken: cfg.Git.Token,
|
GitToken: cfg.Git.Token,
|
||||||
Live: live,
|
Live: live,
|
||||||
|
Notify: a, // авто-уведомления владельцу задачи через Router
|
||||||
}
|
}
|
||||||
a.Router = router
|
a.Router = router
|
||||||
a.Worker = w
|
a.Worker = w
|
||||||
@@ -305,6 +327,16 @@ func (a *App) send(ctx context.Context, uid chat.UserID, text string) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Notify реализует worker.Notifier: авто-уведомление владельцу задачи через
|
||||||
|
// chat.Router.Send (переходы статусов и хендоффы dev↔reviewer со стороны воркера).
|
||||||
|
// Router nil (тесты без Router / ранняя инициализация) — тихо пропускаем.
|
||||||
|
func (a *App) Notify(ctx context.Context, taskID int64, chatID, text string) error {
|
||||||
|
if a.Router == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return a.Router.Send(ctx, chat.UserID(chatID), chat.Message{Text: text})
|
||||||
|
}
|
||||||
|
|
||||||
// cmdName извлекает команду (первое слово до пробела, нижний регистр).
|
// cmdName извлекает команду (первое слово до пробела, нижний регистр).
|
||||||
func cmdName(text string) string {
|
func cmdName(text string) string {
|
||||||
s := strings.TrimSpace(text)
|
s := strings.TrimSpace(text)
|
||||||
@@ -391,23 +423,26 @@ func (a *App) handleBinaryCommand(ctx context.Context, uid chat.UserID, text str
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// helpText — текст команды /help.
|
// helpTextFor — полный текст команды /help с номером версии приложения.
|
||||||
const helpText = "⚡ **Ratatoskr** — помогу продумать задачу и запущу агента.\n\n" +
|
func helpTextFor(version string) string {
|
||||||
"**Команды:**\n" +
|
return "⚡ **Ratatoskr** — помогу продумать задачу и запущу агента.\n\n" +
|
||||||
"`/start` — новая задача\n" +
|
"**Команды:**\n" +
|
||||||
"`/cancel` — отменить текущую\n" +
|
"`/start` — новая задача\n" +
|
||||||
"`/skip` — хватит вопросов, предложить черновик\n" +
|
"`/cancel` — отменить текущую\n" +
|
||||||
"`/retry N` — перезапустить задачу `N`\n" +
|
"`/skip` — хватит вопросов, предложить черновик\n" +
|
||||||
"`/status N` — статус задачи `N`\n" +
|
"`/retry N` — перезапустить задачу `N`\n" +
|
||||||
"`/continue N` — продолжить задачу `N`\n" +
|
"`/status N` — статус задачи `N`\n" +
|
||||||
"`/status` — статус бинаря и обновления\n" +
|
"`/continue N` — продолжить задачу `N`\n" +
|
||||||
"`/update` — применить обновление\n" +
|
"`/status` — статус бинаря и обновления\n" +
|
||||||
"`/help` — эта справка\n\n" +
|
"`/update` — применить обновление\n" +
|
||||||
"Просто опишите задачу — я помогу её продумать."
|
"`/help` — эта справка\n\n" +
|
||||||
|
"Версия: " + version + "\n\n" +
|
||||||
|
"Просто опишите задачу — я помогу её продумать."
|
||||||
|
}
|
||||||
|
|
||||||
// handleHelpCmd отвечает списком команд.
|
// handleHelpCmd отвечает списком команд и номером версии приложения.
|
||||||
func (a *App) handleHelpCmd(ctx context.Context, uid chat.UserID) {
|
func (a *App) handleHelpCmd(ctx context.Context, uid chat.UserID) {
|
||||||
a.send(ctx, uid, helpText)
|
a.send(ctx, uid, helpTextFor(Version))
|
||||||
}
|
}
|
||||||
|
|
||||||
// handleStatusCmd отвечает текущей версией и результатом последней проверки.
|
// handleStatusCmd отвечает текущей версией и результатом последней проверки.
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"regexp"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
)
|
)
|
||||||
@@ -216,17 +217,31 @@ func TestHelp_IsBinaryCommand(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestHelpText_ContainsCommands проверяет, что справка упоминает все команды.
|
// TestHelpText_ContainsCommands проверяет, что справка упоминает все команды
|
||||||
|
// и номер версии приложения.
|
||||||
func TestHelpText_ContainsCommands(t *testing.T) {
|
func TestHelpText_ContainsCommands(t *testing.T) {
|
||||||
|
text := helpTextFor(Version)
|
||||||
for _, cmd := range []string{
|
for _, cmd := range []string{
|
||||||
"/start", "/cancel", "/skip", "/retry", "/status", "/continue",
|
"/start", "/cancel", "/skip", "/retry", "/status", "/continue",
|
||||||
"/update", "/help",
|
"/update", "/help",
|
||||||
} {
|
} {
|
||||||
if !strings.Contains(helpText, cmd) {
|
if !strings.Contains(text, cmd) {
|
||||||
t.Errorf("helpText не упоминает команду %s", cmd)
|
t.Errorf("helpText не упоминает команду %s", cmd)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if !strings.Contains(helpText, "Ratatoskr") {
|
if !strings.Contains(text, "Ratatoskr") {
|
||||||
t.Errorf("helpText не содержит имени бота")
|
t.Errorf("helpText не содержит имени бота")
|
||||||
}
|
}
|
||||||
|
if !strings.Contains(text, Version) {
|
||||||
|
t.Errorf("helpText не содержит номера версии приложения %s", Version)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestVersion_Semver проверяет, что константа версии приложения объявлена
|
||||||
|
// в формате major.minor.patch (например 0.1.0).
|
||||||
|
func TestVersion_Semver(t *testing.T) {
|
||||||
|
re := regexp.MustCompile(`^\d+\.\d+\.\d+$`)
|
||||||
|
if !re.MatchString(Version) {
|
||||||
|
t.Errorf("Version = %q, ожидался формат major.minor.patch", Version)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
828
internal/app/e2e_test.go
Normal file
828
internal/app/e2e_test.go
Normal file
@@ -0,0 +1,828 @@
|
|||||||
|
package app
|
||||||
|
|
||||||
|
// Интеграционный (сквозной) тест «всё приложение от постановки задачи».
|
||||||
|
//
|
||||||
|
// Закрывает оба слоя конвейера одним прогоном, как в проде:
|
||||||
|
//
|
||||||
|
// handleIncoming (/start) → Core.ProcessTurn → Analyst (opencode) → задача ready
|
||||||
|
// Worker.runTask: dev → reviewer → настоящий git push → success
|
||||||
|
//
|
||||||
|
// Аналитик и воркер делят один и тот же *opencode.Runner (как собирает app.New),
|
||||||
|
// а фейк-скрипт opencode различает агентов по argv (аналитик/dev/reviewer) —
|
||||||
|
// возвращая NDJSON-вердикты нужного формата для каждого.
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"os/exec"
|
||||||
|
"path/filepath"
|
||||||
|
"reflect"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/kamelion/ratatoskr-go/internal/analyst"
|
||||||
|
"github.com/kamelion/ratatoskr-go/internal/chat"
|
||||||
|
"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"
|
||||||
|
)
|
||||||
|
|
||||||
|
// ndjsonText собирает строку NDJSON-события opencode с text-партом:
|
||||||
|
// {"type":"text","part":{"text":"<payload>"}}. payload — строковое
|
||||||
|
// представление JSON-вердикта агента (как это делает реальный opencode).
|
||||||
|
func ndjsonText(t *testing.T, payload string) string {
|
||||||
|
t.Helper()
|
||||||
|
b, err := json.Marshal(payload) // экранирует payload как JSON-строку
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("json.Marshal payload: %v", err)
|
||||||
|
}
|
||||||
|
return `{"type":"text","part":{"text":` + string(b) + `}}`
|
||||||
|
}
|
||||||
|
|
||||||
|
// e2eFakeOpenCode пишет shell-скрипт, имитирующий opencode run.
|
||||||
|
// Различает агента по argv ($3 = имя агента после "--agent").
|
||||||
|
//
|
||||||
|
// analyst → NDJSON c вердиктом propose (черновик с репозиторием calc)
|
||||||
|
// dev → простой NDJSON "done"
|
||||||
|
// reviewer→ NDJSON c {"passed":true} в text-парте
|
||||||
|
func e2eFakeOpenCode(t *testing.T, dir string) string {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
analystNDJSON := ndjsonText(t, `{"phase":"propose","title":"Калькулятор","goal":"Сделать веб-калькулятор","repo":"calc","why":"Нужен для учёта","ac":"Работает + - * /","chat_reply":"Черновик готов."}`)
|
||||||
|
reviewerNDJSON := ndjsonText(t, `{"passed":true,"comments":[]}`)
|
||||||
|
|
||||||
|
// каждый вариант печатаем через printf '%s' с одинарными кавычками:
|
||||||
|
// NDJSON содержит двойные кавычки и бэкслеши, но не одинарные — безопасно.
|
||||||
|
analystLine := "printf '%s\\n' '" + analystNDJSON + "'"
|
||||||
|
reviewerLine := "printf '%s\\n' '" + reviewerNDJSON + "'"
|
||||||
|
|
||||||
|
script := `#!/bin/sh
|
||||||
|
agent="$3"
|
||||||
|
case "$agent" in
|
||||||
|
analyst)
|
||||||
|
` + analystLine + `
|
||||||
|
;;
|
||||||
|
reviewer)
|
||||||
|
` + reviewerLine + `
|
||||||
|
;;
|
||||||
|
dev)
|
||||||
|
printf '%%s\n' '{"type":"text","part":{"text":"done"}}'
|
||||||
|
;;
|
||||||
|
*)
|
||||||
|
printf '%%s\n' '{"type":"text","part":{"text":"unknown agent"}}'
|
||||||
|
;;
|
||||||
|
esac
|
||||||
|
exit 0
|
||||||
|
`
|
||||||
|
bin := filepath.Join(dir, "opencode")
|
||||||
|
if err := os.WriteFile(bin, []byte(script), 0o755); err != nil {
|
||||||
|
t.Fatalf("write e2e fake opencode: %v", err)
|
||||||
|
}
|
||||||
|
return bin
|
||||||
|
}
|
||||||
|
|
||||||
|
// e2eAssemble собирает конвейер вручную (те же связи, что app.New),
|
||||||
|
// но с фейк-бинарём, подменённым на e2eFakeOpenCode. Возвращает App,
|
||||||
|
// каталог worktree и fake-канал (для проверки исходящих).
|
||||||
|
func e2eAssemble(t *testing.T) (*App, string, *e2eChannel) {
|
||||||
|
t.Helper()
|
||||||
|
dir := t.TempDir()
|
||||||
|
worktree := filepath.Join(dir, "worktrees")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
store, err := storage.Open(ctx, filepath.Join(dir, "e2e.db"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("storage.Open: %v", err)
|
||||||
|
}
|
||||||
|
t.Cleanup(func() { store.Close() })
|
||||||
|
|
||||||
|
bin := e2eFakeOpenCode(t, dir)
|
||||||
|
runner := &opencode.Runner{
|
||||||
|
Bin: bin,
|
||||||
|
PollInterval: 20 * time.Millisecond,
|
||||||
|
IdleTimeout: 5 * time.Second,
|
||||||
|
HardTimeout: 30 * time.Second,
|
||||||
|
}
|
||||||
|
|
||||||
|
an := &analyst.Analyst{Runner: runner, Worktree: worktree, Agent: "analyst"}
|
||||||
|
coreCtx := core.New(store, an)
|
||||||
|
|
||||||
|
// Router с fake-каналом: ловим исходящие replies для проверки.
|
||||||
|
fake := &e2eChannel{sent: make([]chat.Message, 0)}
|
||||||
|
a := &App{
|
||||||
|
Store: store,
|
||||||
|
CoreCtx: coreCtx,
|
||||||
|
}
|
||||||
|
a.Router = chat.NewRouter(a.handleIncoming)
|
||||||
|
if err := a.Router.Attach(fake); err != nil {
|
||||||
|
t.Fatalf("Attach fake channel: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
w := &worker.Worker{
|
||||||
|
Store: store,
|
||||||
|
Runner: runner,
|
||||||
|
Worktree: worktree,
|
||||||
|
Agent: "dev",
|
||||||
|
Interval: 30 * time.Millisecond,
|
||||||
|
MaxJobs: 1,
|
||||||
|
Live: opencode.NewLiveRegistry(),
|
||||||
|
Notify: a, // авто-уведомления владельцу через Router (как в app.New)
|
||||||
|
// GitToken зададим пустым: origin в seed-репо локальный (file path),
|
||||||
|
// http.extraHeader не нужен для локального пуша.
|
||||||
|
}
|
||||||
|
a.Worker = w
|
||||||
|
return a, worktree, fake
|
||||||
|
}
|
||||||
|
|
||||||
|
// e2eChannel — минимальный fake-канал для перехвата исходящих
|
||||||
|
// и доставки входящих через роутер (как реальный канал).
|
||||||
|
type e2eChannel struct {
|
||||||
|
onMsg chat.Handler
|
||||||
|
sent []chat.Message
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *e2eChannel) Run(_ context.Context) error { return nil }
|
||||||
|
func (c *e2eChannel) OnMessage(h chat.Handler) {
|
||||||
|
c.onMsg = h
|
||||||
|
}
|
||||||
|
func (c *e2eChannel) Send(_ context.Context, _ chat.Address, m chat.Message) error {
|
||||||
|
c.sent = append(c.sent, m)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
func (c *e2eChannel) Ask(_ context.Context, _ chat.Address, m chat.Message) error {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
func (c *e2eChannel) Close() error { return nil }
|
||||||
|
|
||||||
|
// deliver отправляет входящее сообщение через роутер: ставит маршрут
|
||||||
|
// пользователя и вызывает app.handleIncoming (как в проде).
|
||||||
|
func (c *e2eChannel) deliver(uid chat.UserID, text string) {
|
||||||
|
if c.onMsg != nil {
|
||||||
|
c.onMsg(chat.Incoming{UserID: uid, Address: chat.Address("u://" + string(uid)), Msg: chat.Message{Text: text}, Channel: c})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// seedFakeRepo создаёт реальный git-репозиторий для worktree/<repo>:
|
||||||
|
// bare-origin + клон с начальным коммитом на main, чтобы воркер мог
|
||||||
|
// выполнять fetch origin, создавать ветку и пушить. (Дубликат из worker_test,
|
||||||
|
// вынесен сюда, т.к. интеграционный тест живёт в пакете app, а не worker.)
|
||||||
|
func (a *App) seedFakeRepo(t *testing.T, worktree, repo string) {
|
||||||
|
t.Helper()
|
||||||
|
gitRun := func(dir string, args ...string) string {
|
||||||
|
t.Helper()
|
||||||
|
cmd := exec.Command("git", args...)
|
||||||
|
cmd.Dir = dir
|
||||||
|
cmd.Env = append(os.Environ(), "GIT_AUTHOR_NAME=t", "GIT_AUTHOR_EMAIL=t@t",
|
||||||
|
"GIT_COMMITTER_NAME=t", "GIT_COMMITTER_EMAIL=t@t")
|
||||||
|
out, err := cmd.CombinedOutput()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("git %v in %s: %v\n%s", args, dir, err, out)
|
||||||
|
}
|
||||||
|
return string(out)
|
||||||
|
}
|
||||||
|
base := filepath.Join(worktree, repo)
|
||||||
|
if err := os.MkdirAll(base, 0o755); err != nil {
|
||||||
|
t.Fatalf("seed repo %s: %v", repo, err)
|
||||||
|
}
|
||||||
|
origin := filepath.Join(worktree, repo+"-origin.git")
|
||||||
|
gitRun(worktree, "init", "--bare", origin)
|
||||||
|
gitRun(base, "init")
|
||||||
|
gitRun(base, "remote", "add", "origin", origin)
|
||||||
|
if err := os.WriteFile(filepath.Join(base, "seed.txt"), []byte("seed\n"), 0o644); err != nil {
|
||||||
|
t.Fatalf("write seed: %v", err)
|
||||||
|
}
|
||||||
|
gitRun(base, "add", "seed.txt")
|
||||||
|
gitRun(base, "commit", "-m", "seed")
|
||||||
|
gitRun(base, "branch", "-M", "main")
|
||||||
|
gitRun(base, "push", "-u", "origin", "main")
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestE2EWholeAppFromTaskSetup — положительный сквозной путь «всё приложение»:
|
||||||
|
// постановка задачи через аналитика (до ready), затем воркер (dev→reviewer→push)
|
||||||
|
// до success, и проверка, что feature-ветка реально запушена в origin.
|
||||||
|
func TestE2EWholeAppFromTaskSetup(t *testing.T) {
|
||||||
|
a, worktree, fake := e2eAssemble(t)
|
||||||
|
a.seedFakeRepo(t, worktree, "calc")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
uid := chat.UserID("u1")
|
||||||
|
|
||||||
|
// --- 1. Постановка: /start → задача collecting ---
|
||||||
|
fake.deliver(uid, "/start")
|
||||||
|
task, err := a.Store.GetActiveTaskByChatID(ctx, string(uid))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task after /start: %v", err)
|
||||||
|
}
|
||||||
|
if task.Status != storage.StatusCollecting {
|
||||||
|
t.Errorf("status после /start = %q, want collecting", task.Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- 2. Постановка: ответ пользователя → аналитик (opencode) → ready ---
|
||||||
|
// Фейк-analyst возвращает propose с repo=calc → задача должна стать ready.
|
||||||
|
fake.deliver(uid, "Сделай калькулятор в calc")
|
||||||
|
task, err = a.Store.GetTask(ctx, task.ID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task: %v", err)
|
||||||
|
}
|
||||||
|
if task.Status != storage.StatusReady {
|
||||||
|
t.Fatalf("status после analyst = %q, want ready", task.Status)
|
||||||
|
}
|
||||||
|
if len(task.EffectiveRepos()) == 0 {
|
||||||
|
t.Errorf("repos пусто, want [calc]")
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- 3. Согласие «создавай» → задача одобрена (approved) ---
|
||||||
|
fake.deliver(uid, "создавай")
|
||||||
|
task, err = a.Store.GetTask(ctx, task.ID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task: %v", err)
|
||||||
|
}
|
||||||
|
if task.Status != storage.StatusApproved {
|
||||||
|
t.Errorf("status после создавай = %q, want approved", task.Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- 4. Исполнение: воркер (poll) dev → reviewer → push ---
|
||||||
|
// Запускаем реальный poll-диспетчер и ждём, пока задача доедет до success.
|
||||||
|
wkCtx, wkCancel := context.WithCancel(ctx)
|
||||||
|
defer wkCancel()
|
||||||
|
a.Worker.Start(wkCtx)
|
||||||
|
|
||||||
|
deadline := time.Now().Add(30 * time.Second)
|
||||||
|
var last storage.Status
|
||||||
|
for {
|
||||||
|
task, err = a.Store.GetTask(ctx, task.ID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task: %v", err)
|
||||||
|
}
|
||||||
|
last = task.Status
|
||||||
|
if task.Status == storage.StatusSuccess || task.Status == storage.StatusFailed {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
if time.Now().After(deadline) {
|
||||||
|
t.Fatalf("таймаут ожидания success, последний статус %q", last)
|
||||||
|
}
|
||||||
|
time.Sleep(50 * time.Millisecond)
|
||||||
|
}
|
||||||
|
|
||||||
|
task, err = a.Store.GetTask(ctx, task.ID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task: %v", err)
|
||||||
|
}
|
||||||
|
if task.Status != storage.StatusSuccess {
|
||||||
|
t.Fatalf("status после воркера = %q, want success", task.Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
// ревью вызывается ровно один раз (passed сразу) — 2 трассы: dev + reviewer
|
||||||
|
traces, err := a.Store.GetTraces(ctx, task.ID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get traces: %v", err)
|
||||||
|
}
|
||||||
|
if len(traces) != 2 {
|
||||||
|
t.Fatalf("traces = %d, want 2 (dev + reviewer)", len(traces))
|
||||||
|
}
|
||||||
|
if traces[0].Agent != "dev" || traces[0].Status != storage.TraceSuccess {
|
||||||
|
t.Errorf("trace[0] = %s/%s, want dev/success", traces[0].Agent, traces[0].Status)
|
||||||
|
}
|
||||||
|
if traces[1].Agent != "reviewer" || traces[1].Status != storage.TraceSuccess {
|
||||||
|
t.Errorf("trace[1] = %s/%s, want reviewer/success", traces[1].Agent, traces[1].Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- 5. Проверка настоящего push: feature-ветка есть в origin ---
|
||||||
|
branch := "feat/" + task.TaskTag
|
||||||
|
out, err := exec.Command("git", "-C", filepath.Join(worktree, "calc"),
|
||||||
|
"ls-remote", "origin", "refs/heads/"+branch).CombinedOutput()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("ls-remote origin: %v\n%s", err, out)
|
||||||
|
}
|
||||||
|
if !strings.Contains(string(out), "refs/heads/"+branch) {
|
||||||
|
t.Errorf("feature-ветка %q не найдена в origin:\n%s", branch, out)
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- 6. Бот подтвердил через канал (хотя бы одно исходящее было) ---
|
||||||
|
if len(fake.sent) == 0 {
|
||||||
|
t.Error("нет ни одного исходящего сообщения через Router")
|
||||||
|
}
|
||||||
|
t.Logf("E2E OK: задача #%d %s, ветка %s, traces=%d", task.ID, task.Status, branch, len(traces))
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// TestE2ERetryFromStates — положительный сценарий «перезапуск задачи»:
|
||||||
|
// задача перезапускается /retry N из различных стартовых состояний и после
|
||||||
|
// повторного цикла (аналитик → ready → создавай → воркер) снова доходит до
|
||||||
|
// success с повторным реальным push feature-ветки.
|
||||||
|
//
|
||||||
|
// Матрица стартовых состояний:
|
||||||
|
// collecting / ready / cancelled — реальными командами через канал;
|
||||||
|
// success — реальным полным циклом воркера;
|
||||||
|
// failed / timeout — фикстурой через Store (валидной цепочкой переходов),
|
||||||
|
// т.к. механика их достижения в воркере уже покрыта worker_test.go, а
|
||||||
|
// timeout вообще нельзя сэмулировать без реального зависания (RC=-1
|
||||||
|
// ставит только kill по таймауту).
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
// e2eChainToState приводит задачу к заданному статусу валидной цепочкой
|
||||||
|
// переходов через Store.UpdateTask (фикстура для failed/timeout).
|
||||||
|
func e2eChainToState(t *testing.T, store *storage.Storage, id int64, target storage.Status) {
|
||||||
|
t.Helper()
|
||||||
|
ctx := context.Background()
|
||||||
|
// валидные цепочки draft → ... → target (проверяются UpdateTask'ом)
|
||||||
|
var chain []storage.Status
|
||||||
|
switch target {
|
||||||
|
case storage.StatusFailed:
|
||||||
|
chain = []storage.Status{storage.StatusCollecting, storage.StatusReady, storage.StatusApproved, storage.StatusRunning, storage.StatusFailed}
|
||||||
|
case storage.StatusTimeout:
|
||||||
|
chain = []storage.Status{storage.StatusCollecting, storage.StatusReady, storage.StatusApproved, storage.StatusRunning, storage.StatusTimeout}
|
||||||
|
default:
|
||||||
|
t.Fatalf("e2eChainToState: неподдерживаемый target %q", target)
|
||||||
|
}
|
||||||
|
for _, s := range chain {
|
||||||
|
tk, err := store.GetTask(ctx, id)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("e2eChainToState: get task %d: %v", id, err)
|
||||||
|
}
|
||||||
|
tk.Status = s
|
||||||
|
if err := store.UpdateTask(ctx, tk); err != nil {
|
||||||
|
t.Fatalf("e2eChainToState: %s → %s: %v", tk.Status, s, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// e2eWaitWorker запускает poll-диспетчер воркера и ждёт, пока задача дойдёт
|
||||||
|
// до success или failed (лимит 30s). Возвращает итоговый статус.
|
||||||
|
func e2eWaitWorker(t *testing.T, a *App, id int64) storage.Status {
|
||||||
|
t.Helper()
|
||||||
|
ctx := context.Background()
|
||||||
|
wkCtx, wkCancel := context.WithCancel(ctx)
|
||||||
|
defer wkCancel()
|
||||||
|
a.Worker.Start(wkCtx)
|
||||||
|
|
||||||
|
deadline := time.Now().Add(30 * time.Second)
|
||||||
|
var last storage.Status
|
||||||
|
for {
|
||||||
|
tk, err := a.Store.GetTask(ctx, id)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("e2eWaitWorker: get task %d: %v", id, err)
|
||||||
|
}
|
||||||
|
last = tk.Status
|
||||||
|
if tk.Status == storage.StatusSuccess || tk.Status == storage.StatusFailed {
|
||||||
|
return tk.Status
|
||||||
|
}
|
||||||
|
if time.Now().After(deadline) {
|
||||||
|
t.Fatalf("e2eWaitWorker: таймаут, последний статус %q", last)
|
||||||
|
}
|
||||||
|
time.Sleep(50 * time.Millisecond)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// e2eRetryCommon — общий хвост после /retry N для любого стартового состояния.
|
||||||
|
// Проверяет: collecting → аналитик → ready → создавай → воркер → success,
|
||||||
|
// повторный push ветки, и что тег задачи переиспользован (не сменился).
|
||||||
|
func e2eRetryCommon(t *testing.T, a *App, fake *e2eChannel, worktree, uid string, taskID int64) {
|
||||||
|
t.Helper()
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
// тег до retry — должен остаться прежним после повторного цикла
|
||||||
|
before, err := a.Store.GetTask(ctx, taskID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("e2eRetryCommon: get task %d: %v", taskID, err)
|
||||||
|
}
|
||||||
|
origTag := before.TaskTag
|
||||||
|
baseTraces, err := a.Store.GetTraces(ctx, taskID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("e2eRetryCommon: get traces (before): %v", err)
|
||||||
|
}
|
||||||
|
baseN := len(baseTraces)
|
||||||
|
|
||||||
|
// 1. /retry N → collecting, история очищена
|
||||||
|
fake.deliver(chat.UserID(uid), "/retry "+fmt.Sprintf("%d", taskID))
|
||||||
|
tk, err := a.Store.GetTask(ctx, taskID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("e2eRetryCommon: get task after /retry: %v", err)
|
||||||
|
}
|
||||||
|
if tk.Status != storage.StatusCollecting {
|
||||||
|
t.Errorf("retry: status = %q, want collecting", tk.Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 2. новый текст → аналитик → ready (repos не пуст)
|
||||||
|
fake.deliver(chat.UserID(uid), "Переделай: добавь деление")
|
||||||
|
tk, err = a.Store.GetTask(ctx, taskID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("e2eRetryCommon: get task after analyst: %v", err)
|
||||||
|
}
|
||||||
|
if tk.Status != storage.StatusReady {
|
||||||
|
t.Fatalf("retry: status после аналитика = %q, want ready", tk.Status)
|
||||||
|
}
|
||||||
|
if len(tk.EffectiveRepos()) == 0 {
|
||||||
|
t.Errorf("retry: repos пусто, want [calc]")
|
||||||
|
}
|
||||||
|
|
||||||
|
// 3. «создавай» → approved (одобрено, воркер заберёт)
|
||||||
|
fake.deliver(chat.UserID(uid), "создавай")
|
||||||
|
tk, err = a.Store.GetTask(ctx, taskID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("e2eRetryCommon: get task after создавай: %v", err)
|
||||||
|
}
|
||||||
|
if tk.Status != storage.StatusApproved {
|
||||||
|
t.Errorf("retry: status после создавай = %q, want approved", tk.Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 4. воркер → success
|
||||||
|
if got := e2eWaitWorker(t, a, taskID); got != storage.StatusSuccess {
|
||||||
|
t.Fatalf("retry: итог = %q, want success", got)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 5. проверки после повторного цикла
|
||||||
|
tk, err = a.Store.GetTask(ctx, taskID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("e2eRetryCommon: get task final: %v", err)
|
||||||
|
}
|
||||||
|
if tk.TaskTag != origTag {
|
||||||
|
t.Errorf("retry: task_tag сменился %q → %q, want неизменный", origTag, tk.TaskTag)
|
||||||
|
}
|
||||||
|
branch := "feat/" + tk.TaskTag
|
||||||
|
out, err := exec.Command("git", "-C", filepath.Join(worktree, "calc"),
|
||||||
|
"ls-remote", "origin", "refs/heads/"+branch).CombinedOutput()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("retry: ls-remote origin: %v\n%s", err, out)
|
||||||
|
}
|
||||||
|
if !strings.Contains(string(out), "refs/heads/"+branch) {
|
||||||
|
t.Errorf("retry: feature-ветка %q не найдена в origin после повторного push:\n%s", branch, out)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 6. появились ровно 2 новые трассы (dev + reviewer), обе успешные
|
||||||
|
traces, err := a.Store.GetTraces(ctx, taskID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("e2eRetryCommon: get traces (after): %v", err)
|
||||||
|
}
|
||||||
|
if len(traces) != baseN+2 {
|
||||||
|
t.Fatalf("retry: трасс = %d, want %d (base %d + dev + reviewer)", len(traces), baseN+2, baseN)
|
||||||
|
}
|
||||||
|
if traces[baseN].Agent != "dev" || traces[baseN].Status != storage.TraceSuccess {
|
||||||
|
t.Errorf("retry: trace[%d] = %s/%s, want dev/success", baseN, traces[baseN].Agent, traces[baseN].Status)
|
||||||
|
}
|
||||||
|
if traces[baseN+1].Agent != "reviewer" || traces[baseN+1].Status != storage.TraceSuccess {
|
||||||
|
t.Errorf("retry: trace[%d] = %s/%s, want reviewer/success", baseN+1, traces[baseN+1].Agent, traces[baseN+1].Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 7. бот ответил
|
||||||
|
if len(fake.sent) == 0 {
|
||||||
|
t.Error("retry: нет ни одного исходящего сообщения через Router")
|
||||||
|
}
|
||||||
|
t.Logf("RETRY OK: задача #%d %s (из %s), ветка %s, new traces=%d", taskID, tk.Status, before.Status, branch, len(traces)-baseN)
|
||||||
|
}
|
||||||
|
|
||||||
|
// e2eSetupState приводит задачу к стартовому состоянию S и возвращает её ID.
|
||||||
|
// Для состояний, достижимых командами — реальный путь через канал; для
|
||||||
|
// failed/timeout — фикстура через Store.
|
||||||
|
func e2eSetupState(t *testing.T, a *App, fake *e2eChannel, worktree, uid string, state storage.Status) int64 {
|
||||||
|
t.Helper()
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
switch state {
|
||||||
|
case storage.StatusCollecting:
|
||||||
|
fake.deliver(chat.UserID(uid), "/start")
|
||||||
|
tk, err := a.Store.GetActiveTaskByChatID(ctx, uid)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("setup collecting: %v", err)
|
||||||
|
}
|
||||||
|
return tk.ID
|
||||||
|
|
||||||
|
case storage.StatusReady:
|
||||||
|
fake.deliver(chat.UserID(uid), "/start")
|
||||||
|
fake.deliver(chat.UserID(uid), "Сделай калькулятор в calc")
|
||||||
|
tk, err := a.Store.GetActiveTaskByChatID(ctx, uid)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("setup ready: %v", err)
|
||||||
|
}
|
||||||
|
if tk.Status != storage.StatusReady {
|
||||||
|
t.Fatalf("setup ready: status = %q, want ready", tk.Status)
|
||||||
|
}
|
||||||
|
return tk.ID
|
||||||
|
|
||||||
|
case storage.StatusCancelled:
|
||||||
|
fake.deliver(chat.UserID(uid), "/start")
|
||||||
|
tk, err := a.Store.GetActiveTaskByChatID(ctx, uid)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("setup cancelled: get active: %v", err)
|
||||||
|
}
|
||||||
|
fake.deliver(chat.UserID(uid), "/cancel")
|
||||||
|
tk, err = a.Store.GetTask(ctx, tk.ID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("setup cancelled: get task: %v", err)
|
||||||
|
}
|
||||||
|
if tk.Status != storage.StatusCancelled {
|
||||||
|
t.Fatalf("setup cancelled: status = %q, want cancelled", tk.Status)
|
||||||
|
}
|
||||||
|
return tk.ID
|
||||||
|
|
||||||
|
case storage.StatusSuccess:
|
||||||
|
// реальный полный цикл до success
|
||||||
|
fake.deliver(chat.UserID(uid), "/start")
|
||||||
|
fake.deliver(chat.UserID(uid), "Сделай калькулятор в calc")
|
||||||
|
fake.deliver(chat.UserID(uid), "создавай")
|
||||||
|
tk, err := a.Store.GetActiveTaskByChatID(ctx, uid)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("setup success: get active: %v", err)
|
||||||
|
}
|
||||||
|
if got := e2eWaitWorker(t, a, tk.ID); got != storage.StatusSuccess {
|
||||||
|
t.Fatalf("setup success: итог = %q, want success", got)
|
||||||
|
}
|
||||||
|
return tk.ID
|
||||||
|
|
||||||
|
case storage.StatusFailed, storage.StatusTimeout:
|
||||||
|
// фикстура: создаём задачу и прогоняем валидную цепочку переходов
|
||||||
|
id, err := a.Store.CreateTask(ctx, &storage.Task{ChatID: uid})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("setup %s: create task: %v", state, err)
|
||||||
|
}
|
||||||
|
e2eChainToState(t, a.Store, id, state)
|
||||||
|
return id
|
||||||
|
|
||||||
|
default:
|
||||||
|
t.Fatalf("setup: неподдерживаемое состояние %q", state)
|
||||||
|
return 0
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestE2ERetryFromStates — перезапуск задачи /retry N из различных состояний.
|
||||||
|
//
|
||||||
|
// Позитивный блок (задача реально перезапускается и доходит до success):
|
||||||
|
// collecting, ready, failed, timeout.
|
||||||
|
// Негативный блок (завершённые состояния нельзя перезапускать — только /start):
|
||||||
|
// success, cancelled → /retry отклоняется понятным сообщением, статус не меняется.
|
||||||
|
func TestE2ERetryFromStates(t *testing.T) {
|
||||||
|
a, worktree, fake := e2eAssemble(t)
|
||||||
|
a.seedFakeRepo(t, worktree, "calc")
|
||||||
|
|
||||||
|
// --- Позитивный блок: перезапуск и повторный успешный цикл ---
|
||||||
|
positive := []storage.Status{
|
||||||
|
storage.StatusCollecting,
|
||||||
|
storage.StatusReady,
|
||||||
|
storage.StatusFailed,
|
||||||
|
storage.StatusTimeout,
|
||||||
|
}
|
||||||
|
for _, st := range positive {
|
||||||
|
st := st
|
||||||
|
t.Run("positive/"+string(st), func(t *testing.T) {
|
||||||
|
uid := fmt.Sprintf("retry-pos-%s", st)
|
||||||
|
id := e2eSetupState(t, a, fake, worktree, uid, st)
|
||||||
|
e2eRetryCommon(t, a, fake, worktree, uid, id)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- Негативный блок: завершённые состояния перезапуску не подлежат ---
|
||||||
|
terminal := []storage.Status{
|
||||||
|
storage.StatusSuccess,
|
||||||
|
storage.StatusCancelled,
|
||||||
|
}
|
||||||
|
for _, st := range terminal {
|
||||||
|
st := st
|
||||||
|
t.Run("negative/"+string(st), func(t *testing.T) {
|
||||||
|
uid := fmt.Sprintf("retry-neg-%s", st)
|
||||||
|
id := e2eSetupState(t, a, fake, worktree, uid, st)
|
||||||
|
e2eRetryRejected(t, a, fake, uid, id)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// e2eRetryRejected проверяет, что /retry на завершённую задачу НЕ перезапускает:
|
||||||
|
// статус остаётся терминальным, а пользователю уходит понятное сообщение.
|
||||||
|
func e2eRetryRejected(t *testing.T, a *App, fake *e2eChannel, uid string, taskID int64) {
|
||||||
|
t.Helper()
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
before, err := a.Store.GetTask(ctx, taskID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("e2eRetryRejected: get task %d: %v", taskID, err)
|
||||||
|
}
|
||||||
|
if !storage.IsTerminal(before.Status) {
|
||||||
|
t.Fatalf("e2eRetryRejected: предпосылка — статус %q должен быть терминальным", before.Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
// сбрасываем перехваченные исходящие, чтобы ловить только ответ на /retry
|
||||||
|
fake.sent = nil
|
||||||
|
fake.deliver(chat.UserID(uid), "/retry "+fmt.Sprintf("%d", taskID))
|
||||||
|
|
||||||
|
after, err := a.Store.GetTask(ctx, taskID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("e2eRetryRejected: get task after /retry: %v", err)
|
||||||
|
}
|
||||||
|
if after.Status != before.Status {
|
||||||
|
t.Errorf("retry(neg): статус изменился %q → %q, want без изменений", before.Status, after.Status)
|
||||||
|
}
|
||||||
|
if after.TaskTag != before.TaskTag {
|
||||||
|
t.Errorf("retry(neg): task_tag сменился, want неизменный")
|
||||||
|
}
|
||||||
|
|
||||||
|
// понятное сообщение пользователю: говорим про невозможность перезапуска
|
||||||
|
found := false
|
||||||
|
for _, m := range fake.sent {
|
||||||
|
if strings.Contains(m.Text, "нельзя перезапустить") {
|
||||||
|
found = true
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if !found {
|
||||||
|
t.Errorf("retry(neg): нет понятного ответа «нельзя перезапустить», отправлено: %d сообщений", len(fake.sent))
|
||||||
|
}
|
||||||
|
t.Logf("RETRY-REJECTED OK: задача #%d осталась %q", taskID, after.Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestE2EWorkerDoesNotTakeUnconfirmed — регрессия бага, когда воркер брал
|
||||||
|
// задачу на выполнение ещё до «создавай» (по статусу ready, без одобрения).
|
||||||
|
// Проверяет: пока задача в ready, воркер её НЕ трогает; она уходит в работу
|
||||||
|
// только после «создавай» → approved.
|
||||||
|
func TestE2EWorkerDoesNotTakeUnconfirmed(t *testing.T) {
|
||||||
|
a, worktree, fake := e2eAssemble(t)
|
||||||
|
a.seedFakeRepo(t, worktree, "calc")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
uid := chat.UserID("unconf")
|
||||||
|
|
||||||
|
// --- 1. Постановка до ready, без «создавай» ---
|
||||||
|
fake.deliver(uid, "/start")
|
||||||
|
fake.deliver(uid, "Сделай калькулятор в calc")
|
||||||
|
task, err := a.Store.GetActiveTaskByChatID(ctx, string(uid))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task: %v", err)
|
||||||
|
}
|
||||||
|
if task.Status != storage.StatusReady {
|
||||||
|
t.Fatalf("status после analyst = %q, want ready (черновик готов, но НЕ одобрен)", task.Status)
|
||||||
|
}
|
||||||
|
beforeTraces, _ := a.Store.GetTraces(ctx, task.ID)
|
||||||
|
|
||||||
|
// --- 2. Запускаем воркер и даём ему время «промахнуться» ---
|
||||||
|
wkCtx, wkCancel := context.WithCancel(ctx)
|
||||||
|
defer wkCancel()
|
||||||
|
a.Worker.Start(wkCtx)
|
||||||
|
|
||||||
|
time.Sleep(600 * time.Millisecond) // несколько poll-итераций (Interval=30ms)
|
||||||
|
|
||||||
|
// --- 3. Проверяем, что воркер НЕ взял неодобренную задачу ---
|
||||||
|
after, err := a.Store.GetTask(ctx, task.ID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task after worker: %v", err)
|
||||||
|
}
|
||||||
|
if after.Status != storage.StatusReady {
|
||||||
|
t.Fatalf("БАГ: воркер взял неодобренную задачу — status = %q, want ready. Задача должна ждать «создавай».", after.Status)
|
||||||
|
}
|
||||||
|
afterTraces, _ := a.Store.GetTraces(ctx, task.ID)
|
||||||
|
if len(afterTraces) != len(beforeTraces) {
|
||||||
|
t.Errorf("БАГ: появились трассы без одобрения (было %d, стало %d)", len(beforeTraces), len(afterTraces))
|
||||||
|
}
|
||||||
|
// ветка не должна была создаться
|
||||||
|
branch := "feat/" + after.TaskTag
|
||||||
|
out, _ := exec.Command("git", "-C", filepath.Join(worktree, "calc"),
|
||||||
|
"ls-remote", "origin", "refs/heads/"+branch).CombinedOutput()
|
||||||
|
if strings.Contains(string(out), "refs/heads/"+branch) {
|
||||||
|
t.Errorf("БАГ: feature-ветка %q уже в origin до одобрения", branch)
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- 4. «создавай» → approved, теперь воркер берёт и доезжает до success ---
|
||||||
|
fake.deliver(uid, "создавай")
|
||||||
|
tk, err := a.Store.GetTask(ctx, task.ID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task after создавай: %v", err)
|
||||||
|
}
|
||||||
|
if tk.Status != storage.StatusApproved {
|
||||||
|
t.Fatalf("status после создавай = %q, want approved", tk.Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
deadline := time.Now().Add(30 * time.Second)
|
||||||
|
for {
|
||||||
|
tk, err = a.Store.GetTask(ctx, task.ID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task: %v", err)
|
||||||
|
}
|
||||||
|
if tk.Status == storage.StatusSuccess || tk.Status == storage.StatusFailed {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
if time.Now().After(deadline) {
|
||||||
|
t.Fatalf("таймаут ожидания success, последний статус %q", tk.Status)
|
||||||
|
}
|
||||||
|
time.Sleep(50 * time.Millisecond)
|
||||||
|
}
|
||||||
|
if tk.Status != storage.StatusSuccess {
|
||||||
|
t.Fatalf("status после воркера = %q, want success", tk.Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
// ветка теперь должна быть в origin
|
||||||
|
out, err = exec.Command("git", "-C", filepath.Join(worktree, "calc"),
|
||||||
|
"ls-remote", "origin", "refs/heads/"+branch).CombinedOutput()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("ls-remote origin: %v\n%s", err, out)
|
||||||
|
}
|
||||||
|
if !strings.Contains(string(out), "refs/heads/"+branch) {
|
||||||
|
t.Errorf("feature-ветка %q не найдена в origin после одобрения:\n%s", branch, out)
|
||||||
|
}
|
||||||
|
|
||||||
|
t.Logf("WORKER-CONSENT OK: в ready воркер не трогал, после «создавай» → success, ветка %s", branch)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestE2ENotificationsOnTransitions — авто-уведомления владельцу на каждый
|
||||||
|
// переход статуса задачи и отсутствие задвоения на approved.
|
||||||
|
//
|
||||||
|
// Проверяет сквозной поток сообщений через chat.Router:
|
||||||
|
// - ready → approved («создавай»): ровно ОДНО сообщение (ответ «одобрена»),
|
||||||
|
// отдельное уведомление "Задача #N: approved" НЕ дублируется;
|
||||||
|
// - со стороны воркера: уведомления running → dev→reviewer → success
|
||||||
|
// доходят владельцу через Router.Send.
|
||||||
|
func TestE2ENotificationsOnTransitions(t *testing.T) {
|
||||||
|
a, worktree, fake := e2eAssemble(t)
|
||||||
|
a.seedFakeRepo(t, worktree, "calc")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
uid := chat.UserID("u-notif")
|
||||||
|
|
||||||
|
// --- 1. постановка: /start → draft→collecting (приветствие) ---
|
||||||
|
fake.deliver(uid, "/start")
|
||||||
|
task, err := a.Store.GetActiveTaskByChatID(ctx, string(uid))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task after /start: %v", err)
|
||||||
|
}
|
||||||
|
taskID := task.ID
|
||||||
|
|
||||||
|
// --- 2. текст → collecting→ready (сводка черновика) ---
|
||||||
|
fake.deliver(uid, "Сделай калькулятор в calc")
|
||||||
|
task, err = a.Store.GetTask(ctx, taskID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task: %v", err)
|
||||||
|
}
|
||||||
|
if task.Status != storage.StatusReady {
|
||||||
|
t.Fatalf("status = %q, want ready", task.Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- 3. «создавай» → ready→approved: один ответ, без дубля-уведомления ---
|
||||||
|
fake.deliver(uid, "создавай")
|
||||||
|
task, err = a.Store.GetTask(ctx, taskID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task after создавай: %v", err)
|
||||||
|
}
|
||||||
|
if task.Status != storage.StatusApproved {
|
||||||
|
t.Fatalf("status = %q, want approved", task.Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
approvedReplies := 0
|
||||||
|
var approvedText string
|
||||||
|
approvedNotifs := 0
|
||||||
|
for _, m := range fake.sent {
|
||||||
|
if strings.Contains(m.Text, "одобрена") {
|
||||||
|
approvedReplies++
|
||||||
|
approvedText = m.Text
|
||||||
|
}
|
||||||
|
if strings.Contains(m.Text, fmt.Sprintf("Задача #%d: approved", taskID)) {
|
||||||
|
approvedNotifs++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if approvedReplies != 1 {
|
||||||
|
t.Errorf("сообщений об одобрении = %d, want ровно 1 (нет задвоения)", approvedReplies)
|
||||||
|
}
|
||||||
|
if approvedNotifs != 0 {
|
||||||
|
t.Errorf("отдельное уведомление 'Задача #%d: approved' продублировано (%d раз)", taskID, approvedNotifs)
|
||||||
|
}
|
||||||
|
if !strings.Contains(approvedText, "одобрена") {
|
||||||
|
t.Errorf("ответ об одобрении = %q", approvedText)
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- 4. воркер: running → dev→reviewer → success через Router ---
|
||||||
|
wkCtx, wkCancel := context.WithCancel(ctx)
|
||||||
|
defer wkCancel()
|
||||||
|
a.Worker.Start(wkCtx)
|
||||||
|
|
||||||
|
deadline := time.Now().Add(30 * time.Second)
|
||||||
|
for {
|
||||||
|
tk, err := a.Store.GetTask(ctx, taskID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get task: %v", err)
|
||||||
|
}
|
||||||
|
if tk.Status == storage.StatusSuccess || tk.Status == storage.StatusFailed {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
if time.Now().After(deadline) {
|
||||||
|
t.Fatalf("таймаут ожидания success, последний статус %q", tk.Status)
|
||||||
|
}
|
||||||
|
time.Sleep(50 * time.Millisecond)
|
||||||
|
}
|
||||||
|
|
||||||
|
// финальные статус-уведомления воркера в правильном порядке
|
||||||
|
wantOrder := []string{
|
||||||
|
fmt.Sprintf("Задача #%d: running", taskID),
|
||||||
|
fmt.Sprintf("Задача #%d: dev → reviewer (итерация 1)", taskID),
|
||||||
|
fmt.Sprintf("Задача #%d: success", taskID),
|
||||||
|
}
|
||||||
|
var gotOrder []string
|
||||||
|
for _, m := range fake.sent {
|
||||||
|
if strings.HasPrefix(m.Text, fmt.Sprintf("Задача #%d:", taskID)) {
|
||||||
|
gotOrder = append(gotOrder, m.Text)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if !reflect.DeepEqual(gotOrder, wantOrder) {
|
||||||
|
t.Errorf("уведомления воркера = %#v, want %#v", gotOrder, wantOrder)
|
||||||
|
}
|
||||||
|
|
||||||
|
t.Logf("NOTIFICATIONS OK: задача #%d, сообщений одобрения=%d, уведомления=%v", taskID, approvedReplies, gotOrder)
|
||||||
|
}
|
||||||
@@ -52,7 +52,7 @@ func (c *Core) ProcessTurn(ctx context.Context, taskID int64, text string) (Resu
|
|||||||
return Result{}, err
|
return Result{}, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// 2. согласие в фазе ready → создание
|
// 2. согласие в фазе ready → одобрение
|
||||||
if task.Status == storage.StatusReady {
|
if task.Status == storage.StatusReady {
|
||||||
if isConsent(text) {
|
if isConsent(text) {
|
||||||
return c.handleConsent(ctx, task)
|
return c.handleConsent(ctx, task)
|
||||||
@@ -61,6 +61,17 @@ func (c *Core) ProcessTurn(ctx context.Context, taskID int64, text string) (Resu
|
|||||||
return c.handleEdit(ctx, task, text)
|
return c.handleEdit(ctx, task, text)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 2b. approved — финальное одобрение, правка запрещена.
|
||||||
|
// Воркер уже взял/заберёт задачу; текст не меняет статус.
|
||||||
|
if task.Status == storage.StatusApproved {
|
||||||
|
return Result{
|
||||||
|
Reply: "Задача уже одобрена и передана на выполнение. Следите за статусом: /status " + itoa(task.ID),
|
||||||
|
Action: "send",
|
||||||
|
TaskID: task.ID,
|
||||||
|
Status: task.Status,
|
||||||
|
}, nil
|
||||||
|
}
|
||||||
|
|
||||||
// 3. обычный ход: накопление + аналитик
|
// 3. обычный ход: накопление + аналитик
|
||||||
return c.handleTurn(ctx, task, text)
|
return c.handleTurn(ctx, task, text)
|
||||||
}
|
}
|
||||||
@@ -188,6 +199,16 @@ func (c *Core) handleRetry(ctx context.Context, rest string) (Result, error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return c.notFoundReply(ctx, id, err)
|
return c.notFoundReply(ctx, id, err)
|
||||||
}
|
}
|
||||||
|
// Завершённые задачи (success/cancelled/aborted/closed) перезапускать нельзя:
|
||||||
|
// переход → collecting для них невалиден. Только новая задача через /start.
|
||||||
|
if storage.IsTerminal(task.Status) {
|
||||||
|
return Result{
|
||||||
|
Reply: "Задачу #" + itoa(id) + " нельзя перезапустить — она завершена (" + string(task.Status) + "). Создайте новую через /start.",
|
||||||
|
Action: "send",
|
||||||
|
TaskID: id,
|
||||||
|
Status: task.Status,
|
||||||
|
}, nil
|
||||||
|
}
|
||||||
task.Status = storage.StatusCollecting
|
task.Status = storage.StatusCollecting
|
||||||
if err := c.Store.ClearHistory(ctx, id); err != nil {
|
if err := c.Store.ClearHistory(ctx, id); err != nil {
|
||||||
return Result{}, err
|
return Result{}, err
|
||||||
@@ -261,14 +282,15 @@ func (c *Core) notFoundReply(ctx context.Context, id int64, err error) (Result,
|
|||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// handleConsent создаёт задачу (статус ready → ...). Пока — подтверждение готовности.
|
// handleConsent одобряет задачу: ready → approved (финальное одобрение,
|
||||||
|
// после которого воркер забирает задачу на выполнение).
|
||||||
func (c *Core) handleConsent(ctx context.Context, task *storage.Task) (Result, error) {
|
func (c *Core) handleConsent(ctx context.Context, task *storage.Task) (Result, error) {
|
||||||
task.Status = storage.StatusReady
|
task.Status = storage.StatusApproved
|
||||||
if err := c.Store.UpdateTask(ctx, task); err != nil {
|
if err := c.Store.UpdateTask(ctx, task); err != nil {
|
||||||
return Result{}, err
|
return Result{}, err
|
||||||
}
|
}
|
||||||
return Result{
|
return Result{
|
||||||
Reply: "✅ Задача #" + itoa(task.ID) + " готова к запуску.",
|
Reply: "✅ Задача #" + itoa(task.ID) + " одобрена. Запускаю выполнение.",
|
||||||
Action: "created:" + itoa(task.ID),
|
Action: "created:" + itoa(task.ID),
|
||||||
TaskID: task.ID,
|
TaskID: task.ID,
|
||||||
Status: task.Status,
|
Status: task.Status,
|
||||||
|
|||||||
@@ -160,6 +160,10 @@ func TestConsentInReady(t *testing.T) {
|
|||||||
if res.Action != "created:"+itoa(id) {
|
if res.Action != "created:"+itoa(id) {
|
||||||
t.Fatalf("action = %q, want created:%d", res.Action, id)
|
t.Fatalf("action = %q, want created:%d", res.Action, id)
|
||||||
}
|
}
|
||||||
|
task, _ := store.GetTask(ctx, id)
|
||||||
|
if task.Status != storage.StatusApproved {
|
||||||
|
t.Fatalf("status после создавай = %s, want approved", task.Status)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestEditInReadyGoesCollecting(t *testing.T) {
|
func TestEditInReadyGoesCollecting(t *testing.T) {
|
||||||
|
|||||||
@@ -8,7 +8,8 @@ type Status string
|
|||||||
const (
|
const (
|
||||||
StatusDraft Status = "draft" // только что создана
|
StatusDraft Status = "draft" // только что создана
|
||||||
StatusCollecting Status = "collecting" // аналитик собирает детали
|
StatusCollecting Status = "collecting" // аналитик собирает детали
|
||||||
StatusReady Status = "ready" // черновик готов, ждёт запуска
|
StatusReady Status = "ready" // черновик готов, ждёт одобрения пользователя
|
||||||
|
StatusApproved Status = "approved" // пользователь одобрил («создавай») — воркер берёт в работу
|
||||||
StatusRunning Status = "running" // opencode работает
|
StatusRunning Status = "running" // opencode работает
|
||||||
StatusSuccess Status = "success" // задача выполнена
|
StatusSuccess Status = "success" // задача выполнена
|
||||||
StatusFailed Status = "failed" // ошибка выполнения
|
StatusFailed Status = "failed" // ошибка выполнения
|
||||||
@@ -20,23 +21,26 @@ const (
|
|||||||
|
|
||||||
// AllStatuses — все возможные статусы для валидации.
|
// AllStatuses — все возможные статусы для валидации.
|
||||||
var AllStatuses = []Status{
|
var AllStatuses = []Status{
|
||||||
StatusDraft, StatusCollecting, StatusReady,
|
StatusDraft, StatusCollecting, StatusReady, StatusApproved,
|
||||||
StatusRunning, StatusSuccess, StatusFailed, StatusTimeout,
|
StatusRunning, StatusSuccess, StatusFailed, StatusTimeout,
|
||||||
StatusCancelled, StatusAborted, StatusClosed,
|
StatusCancelled, StatusAborted, StatusClosed,
|
||||||
}
|
}
|
||||||
|
|
||||||
// validTransitions задаёт разрешённые переходы статусов.
|
// validTransitions задаёт разрешённые переходы статусов.
|
||||||
var validTransitions = map[Status][]Status{
|
var validTransitions = map[Status][]Status{
|
||||||
StatusDraft: {StatusCollecting, StatusCancelled, StatusAborted},
|
StatusDraft: {StatusCollecting, StatusCancelled, StatusAborted},
|
||||||
StatusCollecting: {StatusReady, StatusDraft, StatusCancelled, StatusAborted},
|
StatusCollecting: {StatusReady, StatusDraft, StatusCancelled, StatusAborted},
|
||||||
StatusReady: {StatusRunning, StatusCancelled, StatusAborted, StatusClosed, StatusCollecting}, // правка готового
|
// ready — черновик готов: «создавай» → approved, либо правка/отмена/закрытие.
|
||||||
StatusRunning: {StatusSuccess, StatusFailed, StatusTimeout, StatusCancelled},
|
StatusReady: {StatusApproved, StatusCancelled, StatusAborted, StatusClosed, StatusCollecting},
|
||||||
StatusSuccess: {StatusClosed},
|
// approved — финальное одобрение: воркер берёт в running, либо отмена/сбой/закрытие.
|
||||||
StatusFailed: {StatusReady, StatusClosed, StatusCancelled, StatusCollecting}, // retry: перезапуск сбора
|
StatusApproved: {StatusRunning, StatusCancelled, StatusAborted, StatusClosed},
|
||||||
StatusTimeout: {StatusReady, StatusClosed, StatusCancelled, StatusCollecting}, // retry: перезапуск сбора
|
StatusRunning: {StatusSuccess, StatusFailed, StatusTimeout, StatusCancelled},
|
||||||
StatusCancelled: {StatusClosed},
|
StatusSuccess: {StatusClosed},
|
||||||
StatusAborted: {StatusClosed},
|
StatusFailed: {StatusReady, StatusClosed, StatusCancelled, StatusCollecting}, // retry: перезапуск сбора
|
||||||
StatusClosed: {}, // терминальный
|
StatusTimeout: {StatusReady, StatusClosed, StatusCancelled, StatusCollecting}, // retry: перезапуск сбора
|
||||||
|
StatusCancelled: {StatusClosed},
|
||||||
|
StatusAborted: {StatusClosed},
|
||||||
|
StatusClosed: {}, // терминальный
|
||||||
}
|
}
|
||||||
|
|
||||||
// IsValidTransition проверяет, допустим ли переход from → to.
|
// IsValidTransition проверяет, допустим ли переход from → to.
|
||||||
|
|||||||
@@ -25,6 +25,12 @@ func TestIsValidTransition(t *testing.T) {
|
|||||||
{StatusTimeout, StatusReady, true}, // retry
|
{StatusTimeout, StatusReady, true}, // retry
|
||||||
{StatusTimeout, StatusCollecting, true}, // retry: перезапуск сбора
|
{StatusTimeout, StatusCollecting, true}, // retry: перезапуск сбора
|
||||||
{StatusTimeout, StatusRunning, false},
|
{StatusTimeout, StatusRunning, false},
|
||||||
|
{StatusReady, StatusApproved, true}, // «создавай» → одобрено
|
||||||
|
{StatusReady, StatusRunning, false}, // без одобрения воркер не запускает
|
||||||
|
{StatusApproved, StatusRunning, true}, // воркер берёт approved
|
||||||
|
{StatusApproved, StatusCancelled, true}, // отмена до запуска
|
||||||
|
{StatusApproved, StatusReady, false}, // финал: назад нельзя
|
||||||
|
{StatusApproved, StatusCollecting, false},
|
||||||
{StatusClosed, StatusDraft, false},
|
{StatusClosed, StatusDraft, false},
|
||||||
{StatusClosed, StatusRunning, false},
|
{StatusClosed, StatusRunning, false},
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -22,6 +22,14 @@ type OpenCodeRunner interface {
|
|||||||
// PollTaskFunc — callback для обработки готовой задачи (подменяемый в тестах).
|
// PollTaskFunc — callback для обработки готовой задачи (подменяемый в тестах).
|
||||||
type PollTaskFunc func(ctx context.Context) error
|
type PollTaskFunc func(ctx context.Context) error
|
||||||
|
|
||||||
|
// Notifier — механизм отправки авто-уведомлений владельцу задачи во время
|
||||||
|
// выполнения. В проде реализуется *app.App через chat.Router.Send (см.
|
||||||
|
// internal/app/app.go → App.Notify); в тестах worker подменяется фейковым
|
||||||
|
// нотифаером. nil — уведомления выключены (ничего не отправляется).
|
||||||
|
type Notifier interface {
|
||||||
|
Notify(ctx context.Context, taskID int64, chatID, text string) error
|
||||||
|
}
|
||||||
|
|
||||||
// Worker — планировщик, запускающий готовые задачи (status=ready → running → success/failed/timeout).
|
// Worker — планировщик, запускающий готовые задачи (status=ready → running → success/failed/timeout).
|
||||||
type Worker struct {
|
type Worker struct {
|
||||||
Store *storage.Storage
|
Store *storage.Storage
|
||||||
@@ -39,6 +47,10 @@ type Worker struct {
|
|||||||
// Через него Runner пишет live-шаги задачи; nil — наблюдение выключено.
|
// Через него Runner пишет live-шаги задачи; nil — наблюдение выключено.
|
||||||
Live *opencode.LiveRegistry
|
Live *opencode.LiveRegistry
|
||||||
|
|
||||||
|
// Notify — нотифаер авто-уведомлений владельцу задачи (статусы + хендоффы
|
||||||
|
// dev↔reviewer). nil — уведомления выключены.
|
||||||
|
Notify Notifier
|
||||||
|
|
||||||
sem chan struct{} // семафор
|
sem chan struct{} // семафор
|
||||||
cancel context.CancelFunc
|
cancel context.CancelFunc
|
||||||
|
|
||||||
@@ -55,6 +67,27 @@ func (w *Worker) runCtx(ctx context.Context, taskID int64) context.Context {
|
|||||||
return opencode.WithLive(ctx, w.Live, taskID)
|
return opencode.WithLive(ctx, w.Live, taskID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// notify отправляет авто-уведомление владельцу задачи, если нотифаер задан.
|
||||||
|
func (w *Worker) notify(ctx context.Context, task *storage.Task, text string) {
|
||||||
|
if w.Notify == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if err := w.Notify.Notify(ctx, task.ID, task.ChatID, text); err != nil {
|
||||||
|
log.Printf("worker: task %d: уведомление: %v", task.ID, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// notifyStatus — уведомление о смене статуса задачи (номер задачи + статус).
|
||||||
|
func (w *Worker) notifyStatus(ctx context.Context, task *storage.Task, s storage.Status) {
|
||||||
|
w.notify(ctx, task, fmt.Sprintf("Задача #%d: %s", task.ID, s))
|
||||||
|
}
|
||||||
|
|
||||||
|
// notifyHandoff — уведомление о передаче задачи между агентами конвейера
|
||||||
|
// на заданной итерации (1-based).
|
||||||
|
func (w *Worker) notifyHandoff(ctx context.Context, task *storage.Task, from, to string, iteration int) {
|
||||||
|
w.notify(ctx, task, fmt.Sprintf("Задача #%d: %s → %s (итерация %d)", task.ID, from, to, iteration))
|
||||||
|
}
|
||||||
|
|
||||||
// Start запускает цикл опроса в фоновой горутине.
|
// Start запускает цикл опроса в фоновой горутине.
|
||||||
func (w *Worker) Start(ctx context.Context) {
|
func (w *Worker) Start(ctx context.Context) {
|
||||||
if w.Agent == "" {
|
if w.Agent == "" {
|
||||||
@@ -113,7 +146,7 @@ func (w *Worker) pollAndDispatch(ctx context.Context) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
tasks, err := w.Store.ListTasks(ctx, storage.TaskFilter{
|
tasks, err := w.Store.ListTasks(ctx, storage.TaskFilter{
|
||||||
Status: storage.StatusReady,
|
Status: storage.StatusApproved,
|
||||||
Limit: slots,
|
Limit: slots,
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -142,7 +175,7 @@ func (w *Worker) pollAndDispatch(ctx context.Context) error {
|
|||||||
// dev дорабатывает по комментариям; прошло → push ветки + success.
|
// dev дорабатывает по комментариям; прошло → push ветки + success.
|
||||||
func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
||||||
// 1. проверяем статус
|
// 1. проверяем статус
|
||||||
if task.Status != storage.StatusReady {
|
if task.Status != storage.StatusApproved {
|
||||||
return fmt.Errorf("%w: task %d status=%q", ErrLaunch, task.ID, task.Status)
|
return fmt.Errorf("%w: task %d status=%q", ErrLaunch, task.ID, task.Status)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -163,6 +196,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
|||||||
if err := w.Store.UpdateTask(ctx, task); err != nil {
|
if err := w.Store.UpdateTask(ctx, task); err != nil {
|
||||||
return fmt.Errorf("%w: set running: %v", ErrUpdate, err)
|
return fmt.Errorf("%w: set running: %v", ErrUpdate, err)
|
||||||
}
|
}
|
||||||
|
w.notifyStatus(ctx, task, storage.StatusRunning)
|
||||||
|
|
||||||
// 2b. клонируем недостающие репозитории в общий каталог.
|
// 2b. клонируем недостающие репозитории в общий каталог.
|
||||||
if err := w.prepareRepos(ctx, repos); err != nil {
|
if err := w.prepareRepos(ctx, repos); err != nil {
|
||||||
@@ -231,6 +265,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
|||||||
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
||||||
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
||||||
}
|
}
|
||||||
|
w.notifyStatus(ctx, task, storage.StatusTimeout)
|
||||||
w.finalizeTrace(ctx, traceID, storage.TraceTimeout, output)
|
w.finalizeTrace(ctx, traceID, storage.TraceTimeout, output)
|
||||||
return nil
|
return nil
|
||||||
default:
|
default:
|
||||||
@@ -238,6 +273,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
|||||||
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
||||||
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
||||||
}
|
}
|
||||||
|
w.notifyStatus(ctx, task, storage.StatusFailed)
|
||||||
w.finalizeTrace(ctx, traceID, storage.TraceFailed, output)
|
w.finalizeTrace(ctx, traceID, storage.TraceFailed, output)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
@@ -245,6 +281,9 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
|||||||
// dev завершился RC=0 → сохраняем успех трассы dev.
|
// dev завершился RC=0 → сохраняем успех трассы dev.
|
||||||
w.finalizeTrace(ctx, traceID, storage.TraceSuccess, output)
|
w.finalizeTrace(ctx, traceID, storage.TraceSuccess, output)
|
||||||
|
|
||||||
|
// уведомляем пользователя о передаче dev → reviewer на ревью.
|
||||||
|
w.notifyHandoff(ctx, task, "dev", "reviewer", iter+1)
|
||||||
|
|
||||||
// 8. РЕВЬЮ: собираем diff всей ветки, запускаем reviewer.
|
// 8. РЕВЬЮ: собираем diff всей ветки, запускаем reviewer.
|
||||||
diffText, dErr := w.branchDiffAll(ctx, repos, branch)
|
diffText, dErr := w.branchDiffAll(ctx, repos, branch)
|
||||||
if dErr != nil {
|
if dErr != nil {
|
||||||
@@ -275,6 +314,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
|||||||
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
||||||
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
||||||
}
|
}
|
||||||
|
w.notifyStatus(ctx, task, storage.StatusFailed)
|
||||||
w.finalizeTrace(ctx, reviewTraceID, storage.TraceFailed, reviewOutput+"\n"+explain)
|
w.finalizeTrace(ctx, reviewTraceID, storage.TraceFailed, reviewOutput+"\n"+explain)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
@@ -289,11 +329,13 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
|||||||
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
||||||
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
||||||
}
|
}
|
||||||
|
w.notifyStatus(ctx, task, storage.StatusSuccess)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Не пройдено: если есть итерации — dev дорабатывает.
|
// Не пройдено: если есть итерации — dev дорабатывает.
|
||||||
if iter+1 < maxReviewIterations {
|
if iter+1 < maxReviewIterations {
|
||||||
|
w.notify(ctx, task, fmt.Sprintf("Задача #%d: reviewer → dev на доработку (итерация %d)", task.ID, iter+1))
|
||||||
feedback = verdict.Comments
|
feedback = verdict.Comments
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
@@ -303,6 +345,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
|||||||
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
||||||
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
||||||
}
|
}
|
||||||
|
w.notify(ctx, task, fmt.Sprintf("Задача #%d: failed — ревью не пройдено за %d итераций", task.ID, maxReviewIterations))
|
||||||
explain := fmt.Sprintf("Ревью не пройдено за %d итераций.", maxReviewIterations)
|
explain := fmt.Sprintf("Ревью не пройдено за %d итераций.", maxReviewIterations)
|
||||||
final := reviewOutput + "\n" + explain
|
final := reviewOutput + "\n" + explain
|
||||||
if e := w.Store.UpdateTraceOutput(ctx, reviewTraceID, final); e != nil {
|
if e := w.Store.UpdateTraceOutput(ctx, reviewTraceID, final); e != nil {
|
||||||
@@ -330,12 +373,13 @@ func (w *Worker) reviewWithRetry(ctx context.Context, taskID int64, cwd, prompt
|
|||||||
return v2, out2, tid2, nil
|
return v2, out2, tid2, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// failTask помечает задачу failed.
|
// failTask помечает задачу failed и уведомляет владельца.
|
||||||
func (w *Worker) failTask(ctx context.Context, task *storage.Task) {
|
func (w *Worker) failTask(ctx context.Context, task *storage.Task) {
|
||||||
task.Status = storage.StatusFailed
|
task.Status = storage.StatusFailed
|
||||||
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
||||||
log.Printf("worker: task %d: set failed: %v", task.ID, e)
|
log.Printf("worker: task %d: set failed: %v", task.ID, e)
|
||||||
}
|
}
|
||||||
|
w.notifyStatus(ctx, task, storage.StatusFailed)
|
||||||
}
|
}
|
||||||
|
|
||||||
// finalizeTrace обновляет output и статус трассы.
|
// finalizeTrace обновляет output и статус трассы.
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ import (
|
|||||||
"os"
|
"os"
|
||||||
"os/exec"
|
"os/exec"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"reflect"
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
@@ -17,6 +18,32 @@ import (
|
|||||||
"github.com/kamelion/ratatoskr-go/internal/storage"
|
"github.com/kamelion/ratatoskr-go/internal/storage"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// fakeNotifier — фейковый нотифаер, собирающий все авто-уведомления воркера.
|
||||||
|
type fakeNotifier struct {
|
||||||
|
notifs []notifCall
|
||||||
|
}
|
||||||
|
|
||||||
|
// notifCall — одно перехваченное уведомление.
|
||||||
|
type notifCall struct {
|
||||||
|
taskID int64
|
||||||
|
chatID string
|
||||||
|
text string
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *fakeNotifier) Notify(_ context.Context, taskID int64, chatID, text string) error {
|
||||||
|
f.notifs = append(f.notifs, notifCall{taskID: taskID, chatID: chatID, text: text})
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// notifTexts возвращает тексты уведомлений в порядке отправки.
|
||||||
|
func notifTexts(n *fakeNotifier) []string {
|
||||||
|
texts := make([]string, len(n.notifs))
|
||||||
|
for i, c := range n.notifs {
|
||||||
|
texts[i] = c.text
|
||||||
|
}
|
||||||
|
return texts
|
||||||
|
}
|
||||||
|
|
||||||
type mockRunnerWorker struct {
|
type mockRunnerWorker struct {
|
||||||
result *opencode.Result
|
result *opencode.Result
|
||||||
err error
|
err error
|
||||||
@@ -95,6 +122,10 @@ func createReadyTask(t *testing.T, s *storage.Storage, title string) *storage.Ta
|
|||||||
if err := s.UpdateTask(ctx, task); err != nil {
|
if err := s.UpdateTask(ctx, task); err != nil {
|
||||||
t.Fatalf("set ready: %v", err)
|
t.Fatalf("set ready: %v", err)
|
||||||
}
|
}
|
||||||
|
task.Status = storage.StatusApproved
|
||||||
|
if err := s.UpdateTask(ctx, task); err != nil {
|
||||||
|
t.Fatalf("set approved: %v", err)
|
||||||
|
}
|
||||||
task, _ = s.GetTask(ctx, id)
|
task, _ = s.GetTask(ctx, id)
|
||||||
return task
|
return task
|
||||||
}
|
}
|
||||||
@@ -259,6 +290,186 @@ func TestWorkerReviewMaxIterations(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestWorkerStatusNotifications — happy path: владелец получает уведомления
|
||||||
|
// на каждый переход статуса со стороны воркера (running → success)
|
||||||
|
// и на хендофф dev→reviewer.
|
||||||
|
func TestWorkerStatusNotifications(t *testing.T) {
|
||||||
|
s := setupWorkerDB(t)
|
||||||
|
task := createReadyTask(t, s, "notif-ok")
|
||||||
|
|
||||||
|
n := &fakeNotifier{}
|
||||||
|
w := &Worker{
|
||||||
|
Store: s,
|
||||||
|
Runner: &mockRunnerWorker{result: &opencode.Result{RC: 0, Stdout: "done", SessionID: "sess-1"}},
|
||||||
|
Worktree: t.TempDir(),
|
||||||
|
Agent: "dev",
|
||||||
|
Notify: n,
|
||||||
|
}
|
||||||
|
seedFakeRepo(t, w.Worktree, "notif-ok")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
if err := w.runTask(ctx, task); err != nil {
|
||||||
|
t.Fatalf("runTask: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(n.notifs) != 3 {
|
||||||
|
t.Fatalf("уведомлений = %d, want 3 (running, dev→reviewer, success)", len(n.notifs))
|
||||||
|
}
|
||||||
|
prefix := "Задача #" + strconv.FormatInt(task.ID, 10)
|
||||||
|
want := []string{
|
||||||
|
prefix + ": running",
|
||||||
|
prefix + ": dev → reviewer (итерация 1)",
|
||||||
|
prefix + ": success",
|
||||||
|
}
|
||||||
|
if got := notifTexts(n); !reflect.DeepEqual(got, want) {
|
||||||
|
t.Errorf("уведомления = %#v, want %#v", got, want)
|
||||||
|
}
|
||||||
|
// все уведомления уходят владельцу задачи (task.ChatID)
|
||||||
|
for _, c := range n.notifs {
|
||||||
|
if c.chatID != task.ChatID {
|
||||||
|
t.Errorf("уведомление ушло в %q, want %q", c.chatID, task.ChatID)
|
||||||
|
}
|
||||||
|
if c.taskID != task.ID {
|
||||||
|
t.Errorf("уведомление для задачи %d, want %d", c.taskID, task.ID)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestWorkerHandoffNotifications — цикл dev↔review: уведомления на оба хендоффа
|
||||||
|
// (dev→reviewer и reviewer→dev на доработку) с номером задачи и итерации.
|
||||||
|
func TestWorkerHandoffNotifications(t *testing.T) {
|
||||||
|
s := setupWorkerDB(t)
|
||||||
|
task := createReadyTask(t, s, "notif-loop")
|
||||||
|
|
||||||
|
n := &fakeNotifier{}
|
||||||
|
w := &Worker{
|
||||||
|
Store: s,
|
||||||
|
Runner: &mockRunnerWorker{
|
||||||
|
result: &opencode.Result{RC: 0, Stdout: "done", SessionID: "sess-1"},
|
||||||
|
reviewSequence: []*opencode.Result{
|
||||||
|
reviewFailedRunner(),
|
||||||
|
{RC: 0, Stdout: `{"passed":true,"comments":[]}`},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
Worktree: t.TempDir(),
|
||||||
|
Agent: "dev",
|
||||||
|
Notify: n,
|
||||||
|
}
|
||||||
|
seedFakeRepo(t, w.Worktree, "notif-loop")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
if err := w.runTask(ctx, task); err != nil {
|
||||||
|
t.Fatalf("runTask: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
prefix := "Задача #" + strconv.FormatInt(task.ID, 10)
|
||||||
|
want := []string{
|
||||||
|
prefix + ": running",
|
||||||
|
prefix + ": dev → reviewer (итерация 1)",
|
||||||
|
prefix + ": reviewer → dev на доработку (итерация 1)",
|
||||||
|
prefix + ": dev → reviewer (итерация 2)",
|
||||||
|
prefix + ": success",
|
||||||
|
}
|
||||||
|
if got := notifTexts(n); !reflect.DeepEqual(got, want) {
|
||||||
|
t.Errorf("уведомления = %#v, want %#v", got, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestWorkerIterationsLimitNotification — при исчерпании лимита итераций
|
||||||
|
// владельцу уходит одно уведомление о failed (без дублей с running→failed).
|
||||||
|
func TestWorkerIterationsLimitNotification(t *testing.T) {
|
||||||
|
s := setupWorkerDB(t)
|
||||||
|
task := createReadyTask(t, s, "notif-lim")
|
||||||
|
|
||||||
|
n := &fakeNotifier{}
|
||||||
|
w := &Worker{
|
||||||
|
Store: s,
|
||||||
|
Runner: &mockRunnerWorker{
|
||||||
|
result: &opencode.Result{RC: 0, Stdout: "done", SessionID: "sess-1"},
|
||||||
|
reviewResult: reviewFailedRunner(),
|
||||||
|
},
|
||||||
|
Worktree: t.TempDir(),
|
||||||
|
Agent: "dev",
|
||||||
|
Notify: n,
|
||||||
|
}
|
||||||
|
seedFakeRepo(t, w.Worktree, "notif-lim")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
if err := w.runTask(ctx, task); err != nil {
|
||||||
|
t.Fatalf("runTask: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
texts := notifTexts(n)
|
||||||
|
if len(texts) == 0 {
|
||||||
|
t.Fatal("нет уведомлений")
|
||||||
|
}
|
||||||
|
last := texts[len(texts)-1]
|
||||||
|
if !strings.Contains(last, "failed") {
|
||||||
|
t.Errorf("последнее уведомление = %q, want упоминание failed", last)
|
||||||
|
}
|
||||||
|
if !strings.Contains(last, "итераци") {
|
||||||
|
t.Errorf("последнее уведомление = %q, want упоминание лимита итераций", last)
|
||||||
|
}
|
||||||
|
// ровно одно уведомление о failed (running→failed не задваивается)
|
||||||
|
var failedCount int
|
||||||
|
for _, txt := range texts {
|
||||||
|
if strings.Contains(txt, ": failed") {
|
||||||
|
failedCount++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if failedCount != 1 {
|
||||||
|
t.Errorf("уведомлений о failed = %d, want ровно 1: %#v", failedCount, texts)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestWorkerTimeoutNotification — RC=-1 (таймаут dev) → уведомление о timeout.
|
||||||
|
func TestWorkerTimeoutNotification(t *testing.T) {
|
||||||
|
s := setupWorkerDB(t)
|
||||||
|
task := createReadyTask(t, s, "notif-timeout")
|
||||||
|
|
||||||
|
n := &fakeNotifier{}
|
||||||
|
w := &Worker{
|
||||||
|
Store: s,
|
||||||
|
Runner: &mockRunnerWorker{result: &opencode.Result{RC: -1, Stdout: ""}},
|
||||||
|
Worktree: t.TempDir(),
|
||||||
|
Notify: n,
|
||||||
|
}
|
||||||
|
seedFakeRepo(t, w.Worktree, "notif-timeout")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
_ = w.runTask(ctx, task)
|
||||||
|
|
||||||
|
prefix := "Задача #" + strconv.FormatInt(task.ID, 10)
|
||||||
|
want := []string{prefix + ": running", prefix + ": timeout"}
|
||||||
|
if got := notifTexts(n); !reflect.DeepEqual(got, want) {
|
||||||
|
t.Errorf("уведомления = %#v, want %#v", got, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestWorkerSpawnErrorNotification — сбой запуска dev → уведомление о failed.
|
||||||
|
func TestWorkerSpawnErrorNotification(t *testing.T) {
|
||||||
|
s := setupWorkerDB(t)
|
||||||
|
task := createReadyTask(t, s, "notif-spawn")
|
||||||
|
|
||||||
|
n := &fakeNotifier{}
|
||||||
|
w := &Worker{
|
||||||
|
Store: s,
|
||||||
|
Runner: &mockRunnerWorker{err: errors.New("opencode not found")},
|
||||||
|
Worktree: t.TempDir(),
|
||||||
|
Notify: n,
|
||||||
|
}
|
||||||
|
seedFakeRepo(t, w.Worktree, "notif-spawn")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
_ = w.runTask(ctx, task)
|
||||||
|
|
||||||
|
prefix := "Задача #" + strconv.FormatInt(task.ID, 10)
|
||||||
|
want := []string{prefix + ": running", prefix + ": failed"}
|
||||||
|
if got := notifTexts(n); !reflect.DeepEqual(got, want) {
|
||||||
|
t.Errorf("уведомления = %#v, want %#v", got, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// reviewNDJSONRunner возвращает вердикт ревьюера как реальный NDJSON-поток opencode,
|
// reviewNDJSONRunner возвращает вердикт ревьюера как реальный NDJSON-поток opencode,
|
||||||
// где JSON находится внутри последнего text-парта.
|
// где JSON находится внутри последнего text-парта.
|
||||||
func reviewNDJSONRunner(v *reviewVerdict) *opencode.Result {
|
func reviewNDJSONRunner(v *reviewVerdict) *opencode.Result {
|
||||||
@@ -625,14 +836,14 @@ func TestWorkerSemaphore(t *testing.T) {
|
|||||||
w.pollAndDispatch(ctx)
|
w.pollAndDispatch(ctx)
|
||||||
time.Sleep(200 * time.Millisecond)
|
time.Sleep(200 * time.Millisecond)
|
||||||
|
|
||||||
// 1 должна быть success, 1 — всё ещё ready
|
// 1 должна быть success, 1 — всё ещё approved
|
||||||
success, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess})
|
success, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess})
|
||||||
ready, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusReady})
|
approved, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusApproved})
|
||||||
if len(success) != 1 {
|
if len(success) != 1 {
|
||||||
t.Errorf("success = %d, want 1 (ready=%d)", len(success), len(ready))
|
t.Errorf("success = %d, want 1 (approved=%d)", len(success), len(approved))
|
||||||
}
|
}
|
||||||
if len(ready) != 1 {
|
if len(approved) != 1 {
|
||||||
t.Errorf("ready = %d, want 1", len(ready))
|
t.Errorf("approved = %d, want 1", len(approved))
|
||||||
}
|
}
|
||||||
|
|
||||||
// первая завершилась и вернула токен в сем — можем диспатчить вторую
|
// первая завершилась и вернула токен в сем — можем диспатчить вторую
|
||||||
|
|||||||
Reference in New Issue
Block a user