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" "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 — строковое // представление 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(), // GitToken зададим пустым: origin в seed-репо локальный (file path), // http.extraHeader не нужен для локального пуша. } a.Worker = w return a, worktree, fake } // e2eChannel — минимальный fake-канал для перехвата исходящих // и доставки входящих через роутер (как реальный канал). type e2eChannel struct { onMsg func(chat.Incoming) 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/: // 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) }