- E2E (app): e2eFakeAPI переведён на v2 HTTP API opencode (/api/*) с определением агента по тексту промпта; Router получает processed-счётчик и WaitProcessed, e2eChannel.deliver ждёт асинхронную обработку — убирает гонку «запрос сразу после deliver» и коллатеральный 'database is closed'. - app_test: одинарные YAML-кавычки для путей Windows (backslash-escape) + закрытие Store в TestNew/TestNew_RunCtxCancel/TestNew_UpdateWiring. - config_test: абсолютный путь строится с корнем тома (C:\...) и одинарными кавычками YAML. - opencode/server_test: fakeServeBin на Windows — .cmd с ping (#!/bin/sh не исполняется). - worker_test: TestWorkerSemaphore поллит до целевого статуса вместо фиксированных sleep (git на Windows медленнее).
899 lines
36 KiB
Go
899 lines
36 KiB
Go
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, эмулирующий v2 HTTP API
|
||
// (пути с префиксом /api/*, см. Client в internal/opencode). Агент
|
||
// (analyst/dev/reviewer) определяется по тексту промпта на POST
|
||
// /api/session/{id}/prompt; вердикт возвращается как text-часть завершённого
|
||
// assistant-сообщения, которое отдаёт GET /api/session/{id}/message.
|
||
func e2eFakeAPI(t *testing.T) string {
|
||
t.Helper()
|
||
var mu sync.Mutex
|
||
sessions := map[string]string{} // id → agent
|
||
|
||
agentOf := func(prompt string) string {
|
||
switch {
|
||
case strings.Contains(prompt, "Ты — аналитик"):
|
||
return "analyst"
|
||
case strings.Contains(prompt, "Ты — dev-агент"):
|
||
return "dev"
|
||
case strings.Contains(prompt, "Ты — ревьюер"):
|
||
return "reviewer"
|
||
default:
|
||
return "unknown"
|
||
}
|
||
}
|
||
|
||
verdictFor := func(agent string) string {
|
||
if v, ok := e2eAgentVerdicts[agent]; ok {
|
||
return v
|
||
}
|
||
return "unknown agent"
|
||
}
|
||
|
||
sessionID := func(path, suffix string) string {
|
||
return strings.TrimSuffix(strings.TrimPrefix(path, "/api/session/"), suffix)
|
||
}
|
||
|
||
assistantMsg := func(id, agent string) map[string]any {
|
||
ts := time.Now().UnixMilli()
|
||
return map[string]any{
|
||
"id": "m-" + id,
|
||
"type": "assistant",
|
||
"content": []map[string]any{{"type": "text", "text": verdictFor(agent)}},
|
||
"model": map[string]any{"providerID": "test", "id": "m"},
|
||
"time": map[string]any{"created": ts, "completed": ts},
|
||
}
|
||
}
|
||
|
||
h := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||
switch {
|
||
case r.Method == http.MethodPost && r.URL.Path == "/api/session":
|
||
// v2 create: {model:{...}} → {data:{id}}
|
||
id := fmt.Sprintf("e2e-%d", len(sessions)+1)
|
||
mu.Lock()
|
||
sessions[id] = ""
|
||
mu.Unlock()
|
||
writeJSON(w, map[string]any{"data": map[string]any{"id": id}})
|
||
|
||
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/prompt"):
|
||
// v2 durable admit: {prompt:{text}} → {data:{id,timeCreated}}
|
||
id := sessionID(r.URL.Path, "/prompt")
|
||
var req struct {
|
||
Prompt struct {
|
||
Text string `json:"text"`
|
||
} `json:"prompt"`
|
||
}
|
||
_ = json.NewDecoder(r.Body).Decode(&req)
|
||
mu.Lock()
|
||
sessions[id] = agentOf(req.Prompt.Text)
|
||
mu.Unlock()
|
||
writeJSON(w, map[string]any{"data": map[string]any{"id": "p-" + id, "timeCreated": time.Now().UnixMilli()}})
|
||
|
||
case r.Method == http.MethodGet && strings.HasSuffix(r.URL.Path, "/message"):
|
||
// v2 поллинг: {data:[Session.Message]} (новейшие первыми).
|
||
id := sessionID(r.URL.Path, "/message")
|
||
mu.Lock()
|
||
agent := sessions[id]
|
||
mu.Unlock()
|
||
writeJSON(w, map[string]any{"data": []map[string]any{assistantMsg(id, agent)}})
|
||
|
||
case r.Method == http.MethodGet && r.URL.Path == "/api/session/active":
|
||
// сессий в активных дренажах нет → ответ завершён.
|
||
writeJSON(w, map[string]any{"data": map[string]any{}})
|
||
|
||
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/interrupt"):
|
||
writeJSON(w, map[string]any{"data": 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)
|
||
fake.router = a.Router
|
||
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
|
||
router *chat.Router
|
||
}
|
||
|
||
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 (как в проде). Так как роутер
|
||
// обрабатывает входящие асинхронно (процессор-горутина), deliver ждёт, пока
|
||
// обработка события завершится, — иначе тесты (сразу читающие состояние БД)
|
||
// гоняются с обработчиком.
|
||
func (c *e2eChannel) deliver(uid chat.UserID, text string) {
|
||
if c.onMsg == nil {
|
||
return
|
||
}
|
||
target := int64(0)
|
||
if c.router != nil {
|
||
target = c.router.Processed() + 1
|
||
}
|
||
c.onMsg(chat.Incoming{UserID: uid, Address: chat.Address("u://" + string(uid)), Msg: chat.Message{Text: text}, Channel: c})
|
||
if c.router != nil && !c.router.WaitProcessed(target) {
|
||
panic("e2e: роутер не обработал входящее за 30s")
|
||
}
|
||
}
|
||
|
||
// 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)
|
||
} |