Files
ratatoskr-go/internal/app/e2e_test.go
Hermes 7eb3a0292c
Some checks failed
CI / test (push) Failing after 32s
CI / build-and-package (amd64, linux) (push) Successful in 38s
CI / build-and-package (amd64, windows) (push) Successful in 38s
fix(opencode): убрать смешение слоёв API — перейти целиком на experimental (/session)
Корень проблемы «не получаем результаты»: клиент смешивал два слоя opencode
serve. CreateSession ходил на /api/session (v2, ждал {data.id}), Verdict — на
/api/session/{id}/message?order=desc и ждал {data:[{type,content}]}, где поле
content[].type/text физически отсутствует, поэтому вердикт никогда не находился
и поллинг уходил в вечный таймаут. Abort и вовсе звал несуществующий /interrupt.

Теперь весь код на experimental-слое, как сверено с sst/opencode (ветка dev):
- CreateSession: POST /session → голая Session, id в .id.
- Send: блокирующий POST /session/{id}/message, тело {parts:[{type:text,text}]},
  вердикт из частей parts[].type=="text" ответа. Это и есть результат — метод
  Verdict и отдельный GET удалены.
- textCount (прогресс): GET /session/{id}/message → голый массив [{info, parts}].
- Abort: POST /session/{id}/abort.

Runner: блокирующий Send запускается в горутине (канал вердикта/ошибки),
параллельно поллим textCount (рост text-частей сбрасывает idle-таймер). При
idle/hard-таймауте или отмене контекста — Abort + cancel() Send-горутины → rc=-1.

Send ходит через отдельный http.Client без жёсткого Timeout (управляется ctx),
чтобы длинная генерация не обрывалась на 30s. Тесты/fakeAPIServer переведены на
экспериментальный формат. Версия → 0.2.2.
2026-08-18 20:49:34 +05:00

853 lines
34 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package app
// Интеграционный (сквозной) тест «всё приложение от постановки задачи».
//
// Закрывает оба слоя конвейера одним прогоном, как в проде:
//
// handleIncoming (/start) → Core.ProcessTurn → Analyst (opencode) → задача ready
// Worker.runTask: dev → reviewer → настоящий git push → success
//
// Аналитик и воркер делят один и тот же *opencode.Runner (как собирает app.New),
// а opencode serve эмулируется фейковым HTTP API-сервером (e2eFakeAPI). Агент
// определяется по title сессии (ratatoskr-analyst / ratatoskr-dev / ratatoskr-reviewer),
// вердикты возвращаются как text-части assistant-сообщений.
import (
"context"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"os"
"os/exec"
"path/filepath"
"reflect"
"strings"
"sync"
"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"
)
// вердикты фейкового агента по имени.
var (
e2eAgentVerdicts = map[string]string{
"analyst": `{"phase":"propose","title":"Калькулятор","goal":"Сделать веб-калькулятор","repo":"calc","why":"Нужен для учёта","ac":"Работает + - * /","chat_reply":"Черновик готов."}`,
"dev": `done`,
"reviewer": `{"passed":true,"comments":[]}`,
}
)
// e2eFakeAPI поднимает фейковый opencode serve experimental HTTP API (пути
// БЕЗ /api) и возвращает URL. По title сессии (ratatoskr-<agent>) определяет
// агента и возвращает его вердикт как text-часть ответа на POST /message.
func e2eFakeAPI(t *testing.T) string {
t.Helper()
var mu sync.Mutex
sessions := map[string]string{} // id → agent
verdictFor := func(agent string) string {
if v, ok := e2eAgentVerdicts[agent]; ok {
return v
}
return "unknown agent"
}
h := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch {
case r.Method == http.MethodPost && r.URL.Path == "/session":
var req struct {
Title string `json:"title"`
}
_ = json.NewDecoder(r.Body).Decode(&req)
agent := strings.TrimPrefix(req.Title, "ratatoskr-")
mu.Lock()
id := fmt.Sprintf("e2e-%d", len(sessions)+1)
sessions[id] = agent
mu.Unlock()
// experimental: голая Session, id напрямую.
writeJSON(w, map[string]any{"id": id, "agent": agent, "model": map[string]any{"id": "m"}})
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/message"):
// блокирующий ответ: вердикт как text-часть.
id := strings.TrimSuffix(strings.TrimPrefix(r.URL.Path, "/session/"), "/message")
mu.Lock()
agent := sessions[id]
mu.Unlock()
writeJSON(w, map[string]any{"info": map[string]any{"role": "assistant"}, "parts": []map[string]any{{"type": "text", "text": verdictFor(agent)}}})
case r.Method == http.MethodGet && strings.HasSuffix(r.URL.Path, "/message"):
// поллинг прогресса: голый массив [{info, parts}].
id := strings.TrimSuffix(strings.TrimPrefix(r.URL.Path, "/session/"), "/message")
mu.Lock()
agent := sessions[id]
mu.Unlock()
writeJSON(w, []map[string]any{{"info": map[string]any{"role": "assistant"}, "parts": []map[string]any{{"type": "text", "text": verdictFor(agent)}}}})
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/abort"):
writeJSON(w, map[string]any{"ok": true})
default:
http.NotFound(w, r)
}
})
srv := httptest.NewServer(h)
t.Cleanup(srv.Close)
return srv.URL
}
func writeJSON(w http.ResponseWriter, v any) {
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(v)
}
// e2eAssemble собирает конвейер вручную (те же связи, что app.New),
// но с фейк-сервером opencode (e2eFakeAPI), зарегистрированным в пуле.
// Возвращает 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() })
fakeURL := e2eFakeAPI(t)
pool := opencode.NewPool(worktree)
pool.RegisterExternal(worktree, fakeURL) // worktree обслуживается фейком
runner := &opencode.Runner{
Pool: pool,
PollInterval: 5 * time.Millisecond,
IdleTimeout: 5 * time.Second,
HardTimeout: 30 * time.Second,
}
t.Cleanup(pool.Close)
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)
}