После завершения задачи со статусом failed/timeout воркер запускает постмортем-анализ (агент postmortem): разбирает сессии dev/reviewer, оценивает причины сбоя и шлёт владельцу уведомление с анализом. Статус задачи не меняет; сбои анализа не влияют на исход.
1093 lines
35 KiB
Go
1093 lines
35 KiB
Go
package worker
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"errors"
|
||
"os"
|
||
"os/exec"
|
||
"path/filepath"
|
||
"reflect"
|
||
"strconv"
|
||
"strings"
|
||
"testing"
|
||
"time"
|
||
|
||
"github.com/kamelion/ratatoskr-go/internal/opencode"
|
||
"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 {
|
||
result *opencode.Result
|
||
err error
|
||
|
||
// reviewResult — результат для reviewer-агента; если nil, при ревью
|
||
// вернётся валидный «passed» verdict.
|
||
reviewResult *opencode.Result
|
||
// reviewSequence — если задан, ревью-запуски берут результаты по порядку
|
||
// (последний вечный). Позволяет смоделировать fail-then-pass.
|
||
reviewSequence []*opencode.Result
|
||
reviewCount int
|
||
|
||
// postMortemResult — результат постмортем-агента; nil → RC:0 с тестовым
|
||
// резюме. postMortemCount — сколько раз постмортем вызывался.
|
||
postMortemResult *opencode.Result
|
||
postMortemCount int
|
||
}
|
||
|
||
func (m *mockRunnerWorker) Run(_ context.Context, _, _, agent, _ string) (*opencode.Result, error) {
|
||
if agent == "reviewer" {
|
||
if len(m.reviewSequence) > 0 {
|
||
i := m.reviewCount
|
||
m.reviewCount++
|
||
if i >= len(m.reviewSequence) {
|
||
return m.reviewSequence[len(m.reviewSequence)-1], m.err
|
||
}
|
||
return m.reviewSequence[i], m.err
|
||
}
|
||
if m.reviewResult != nil {
|
||
return m.reviewResult, m.err
|
||
}
|
||
return &opencode.Result{RC: 0, Stdout: `{"passed":true,"comments":[]}`}, nil
|
||
}
|
||
if agent == postMortemAgent {
|
||
m.postMortemCount++
|
||
if m.postMortemResult != nil {
|
||
return m.postMortemResult, m.err
|
||
}
|
||
return &opencode.Result{RC: 0, Stdout: "Почему так случилось:\n- тестовое резюме\nЧто сделать:\n- исправить"}, nil
|
||
}
|
||
return m.result, m.err
|
||
}
|
||
|
||
// reviewFailedRunner возвращает вердикт not-passed с комментариями.
|
||
func reviewFailedRunner() *opencode.Result {
|
||
return &opencode.Result{RC: 0, Stdout: `{"passed":false,"critical_issues":[],"solid_violations":["DIP: высокая связанность"],"comments":["исправь связанность"]}`}
|
||
}
|
||
|
||
// reviewGarbageRunner — не-парсируемый вывод reviewer.
|
||
func reviewGarbageRunner() *opencode.Result {
|
||
return &opencode.Result{RC: 0, Stdout: "не JSON вовсе"}
|
||
}
|
||
|
||
func setupWorkerDB(t *testing.T) *storage.Storage {
|
||
t.Helper()
|
||
ctx := context.Background()
|
||
f := filepath.Join(t.TempDir(), "test.db")
|
||
s, err := storage.Open(ctx, f)
|
||
if err != nil {
|
||
t.Fatalf("open storage: %v", err)
|
||
}
|
||
t.Cleanup(func() { s.Close() })
|
||
return s
|
||
}
|
||
|
||
func createReadyTask(t *testing.T, s *storage.Storage, title string) *storage.Task {
|
||
t.Helper()
|
||
ctx := context.Background()
|
||
task := &storage.Task{
|
||
ChatID: "tg://worker-test",
|
||
Title: title,
|
||
Goal: "сделать " + title,
|
||
Repos: []string{title},
|
||
Why: "для теста",
|
||
AC: "работает",
|
||
TaskTag: "test-" + title,
|
||
}
|
||
id, err := s.CreateTask(ctx, task)
|
||
if err != nil {
|
||
t.Fatalf("create task: %v", err)
|
||
}
|
||
task.ID = id
|
||
task.Status = storage.StatusCollecting
|
||
if err := s.UpdateTask(ctx, task); err != nil {
|
||
t.Fatalf("set collecting: %v", err)
|
||
}
|
||
task.Status = storage.StatusReady
|
||
if err := s.UpdateTask(ctx, task); err != nil {
|
||
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)
|
||
return task
|
||
}
|
||
|
||
// seedFakeRepo создаёт реальный git-репозиторий для worktree/<repo>:
|
||
// bare-origin + клон с начальным коммитом на main, чтобы воркер мог
|
||
// выполнять fetch origin, создавать ветку и пушить.
|
||
func 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)
|
||
}
|
||
// родительский каталог для bare-origin
|
||
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")
|
||
}
|
||
|
||
func TestWorkerHappyPath(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
task := createReadyTask(t, s, "calc")
|
||
|
||
w := &Worker{
|
||
Store: s,
|
||
Runner: &mockRunnerWorker{result: &opencode.Result{RC: 0, Stdout: "done", SessionID: "sess-1"}},
|
||
Worktree: t.TempDir(),
|
||
Agent: "dev",
|
||
}
|
||
seedFakeRepo(t, w.Worktree, "calc")
|
||
|
||
ctx := context.Background()
|
||
if err := w.runTask(ctx, task); err != nil {
|
||
t.Fatalf("runTask: %v", err)
|
||
}
|
||
|
||
task, err := s.GetTask(ctx, task.ID)
|
||
if err != nil {
|
||
t.Fatalf("get task: %v", err)
|
||
}
|
||
if task.Status != storage.StatusSuccess {
|
||
t.Errorf("status = %q, want success", task.Status)
|
||
}
|
||
|
||
traces, err := s.GetTraces(ctx, task.ID)
|
||
if err != nil {
|
||
t.Fatalf("get traces: %v", err)
|
||
}
|
||
if len(traces) != 2 {
|
||
t.Fatalf("got %d traces, want 2 (dev + reviewer)", len(traces))
|
||
}
|
||
if traces[0].Status != storage.TraceSuccess {
|
||
t.Errorf("trace[0] status = %q, want success", traces[0].Status)
|
||
}
|
||
if traces[0].Agent != "dev" {
|
||
t.Errorf("trace[0] agent = %q, want dev", traces[0].Agent)
|
||
}
|
||
if traces[0].SessionID != "sess-1" {
|
||
t.Errorf("trace[0] session = %q, want sess-1", traces[0].SessionID)
|
||
}
|
||
if traces[1].Agent != "reviewer" {
|
||
t.Errorf("trace[1] agent = %q, want reviewer", traces[1].Agent)
|
||
}
|
||
if traces[1].Status != storage.TraceSuccess {
|
||
t.Errorf("trace[1] status = %q, want success", traces[1].Status)
|
||
}
|
||
}
|
||
|
||
// TestWorkerReviewFailThenSuccess проверяет цикл dev↔review: первый вердикт
|
||
// not-passed → dev дорабатывает, второй passed → успех + пуш.
|
||
func TestWorkerReviewFailThenSuccess(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
task := createReadyTask(t, s, "calc2")
|
||
|
||
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",
|
||
}
|
||
seedFakeRepo(t, w.Worktree, "calc2")
|
||
|
||
ctx := context.Background()
|
||
if err := w.runTask(ctx, task); err != nil {
|
||
t.Fatalf("runTask: %v", err)
|
||
}
|
||
|
||
task, _ = s.GetTask(ctx, task.ID)
|
||
if task.Status != storage.StatusSuccess {
|
||
t.Errorf("status = %q, want success", task.Status)
|
||
}
|
||
|
||
traces, err := s.GetTraces(ctx, task.ID)
|
||
if err != nil {
|
||
t.Fatalf("get traces: %v", err)
|
||
}
|
||
// 2 dev + 2 reviewer (обе итерации)
|
||
if len(traces) != 4 {
|
||
t.Fatalf("got %d traces, want 4 (2 dev + 2 reviewer)", len(traces))
|
||
}
|
||
// первая итерация: dev→reviewer; промпт dev второй итерации должен содержать feedback
|
||
if !strings.Contains(traces[2].Prompt, "РЕВЬЮ НЕ ПРОЙДЕНО") {
|
||
t.Error("второй dev-промпт не содержит feedback от ревьюера")
|
||
}
|
||
if !strings.Contains(traces[2].Prompt, "исправь связанность") {
|
||
t.Error("второй dev-промпт не содержит комментарий ревьюера")
|
||
}
|
||
}
|
||
|
||
// TestWorkerReviewMaxIterations: ревью всё время not-passed → задача failed по R4.
|
||
func TestWorkerReviewMaxIterations(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
task := createReadyTask(t, s, "calc3")
|
||
|
||
w := &Worker{
|
||
Store: s,
|
||
Runner: &mockRunnerWorker{
|
||
result: &opencode.Result{RC: 0, Stdout: "done", SessionID: "sess-1"},
|
||
reviewResult: reviewFailedRunner(),
|
||
},
|
||
Worktree: t.TempDir(),
|
||
Agent: "dev",
|
||
}
|
||
seedFakeRepo(t, w.Worktree, "calc3")
|
||
|
||
// mock всегда not-passed → исчерпаем лимит за 3 итерации (maxReviewIterations=3).
|
||
ctx := context.Background()
|
||
if err := w.runTask(ctx, task); err != nil {
|
||
t.Fatalf("runTask: %v", err)
|
||
}
|
||
task, _ = s.GetTask(ctx, task.ID)
|
||
if task.Status != storage.StatusFailed {
|
||
t.Errorf("status = %q, want failed после лимита итераций", task.Status)
|
||
}
|
||
}
|
||
|
||
// 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("нет уведомлений")
|
||
}
|
||
// где-то в списке есть уведомление о failed с упоминанием лимита итераций
|
||
// (над ним — уведомление postmortem, поэтому последним оно не обязано быть).
|
||
var limitNotif string
|
||
for _, txt := range texts {
|
||
if strings.Contains(txt, "failed") && strings.Contains(txt, "итераци") {
|
||
limitNotif = txt
|
||
break
|
||
}
|
||
}
|
||
if limitNotif == "" {
|
||
t.Errorf("нет уведомления о лимите итераций (failed+итераци): %#v", texts)
|
||
}
|
||
// ровно одно уведомление о 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)
|
||
}
|
||
// есть постмортем-уведомление с анализом
|
||
var pmCount int
|
||
for _, txt := range texts {
|
||
if strings.Contains(txt, "анализ (после failed)") {
|
||
pmCount++
|
||
}
|
||
}
|
||
if pmCount != 1 {
|
||
t.Errorf("постмортем-уведомлений = %d, want ровно 1: %#v", pmCount, 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)
|
||
texts := notifTexts(n)
|
||
if len(texts) != 3 {
|
||
t.Fatalf("уведомлений = %d, want 3: %#v", len(texts), texts)
|
||
}
|
||
want := []string{prefix + ": running", prefix + ": timeout"}
|
||
if !reflect.DeepEqual(texts[:2], want) {
|
||
t.Errorf("первые уведомления = %#v, want %#v", texts[:2], want)
|
||
}
|
||
if !strings.Contains(texts[2], "🔍 анализ (после timeout)") {
|
||
t.Errorf("постмортем-уведомление = %q, want упоминание «анализ (после timeout)»", texts[2])
|
||
}
|
||
}
|
||
|
||
// 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)
|
||
texts := notifTexts(n)
|
||
if len(texts) != 3 {
|
||
t.Fatalf("уведомлений = %d, want 3: %#v", len(texts), texts)
|
||
}
|
||
want := []string{prefix + ": running", prefix + ": failed"}
|
||
if !reflect.DeepEqual(texts[:2], want) {
|
||
t.Errorf("первые уведомления = %#v, want %#v", texts[:2], want)
|
||
}
|
||
if !strings.Contains(texts[2], "анализ (после failed)") {
|
||
t.Errorf("постмортем-уведомление = %q, want упоминание «анализ (после failed)»", texts[2])
|
||
}
|
||
}
|
||
|
||
// reviewNDJSONRunner возвращает вердикт ревьюера как реальный NDJSON-поток opencode,
|
||
// где JSON находится внутри последнего text-парта.
|
||
func reviewNDJSONRunner(v *reviewVerdict) *opencode.Result {
|
||
inner, err := json.Marshal(v)
|
||
if err != nil {
|
||
panic(err)
|
||
}
|
||
// text-парт содержит вердикт как JSON-строку (экранированную внутри NDJSON).
|
||
return &opencode.Result{RC: 0, Stdout: "{\"type\":\"step_start\"}\n{\"type\":\"text\",\"part\":{\"text\":" + strconv.Quote(string(inner)) + "}}"}
|
||
}
|
||
|
||
// TestWorkerReviewNDJSONPass: ревьюер вернул passed:true внутри NDJSON-потока —
|
||
// извлечение должно достать вердикт, а не зафейлить задачу (регрессия двойного ревью).
|
||
func TestWorkerReviewNDJSONPass(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
task := createReadyTask(t, s, "calcnd")
|
||
|
||
w := &Worker{
|
||
Store: s,
|
||
Runner: &mockRunnerWorker{
|
||
result: &opencode.Result{RC: 0, Stdout: "done", SessionID: "sess-1"},
|
||
reviewResult: reviewNDJSONRunner(&reviewVerdict{
|
||
Passed: true,
|
||
Comments: []string{"ок"},
|
||
}),
|
||
},
|
||
Worktree: t.TempDir(),
|
||
Agent: "dev",
|
||
}
|
||
seedFakeRepo(t, w.Worktree, "calcnd")
|
||
|
||
ctx := context.Background()
|
||
if err := w.runTask(ctx, task); err != nil {
|
||
t.Fatalf("runTask: %v", err)
|
||
}
|
||
task, _ = s.GetTask(ctx, task.ID)
|
||
if task.Status != storage.StatusSuccess {
|
||
t.Errorf("status = %q, want success (вердикт из NDJSON должен распарситься)", task.Status)
|
||
}
|
||
|
||
// ревью должен был запуститься ровно 1 раз (никакого retry на «невалидный»).
|
||
traces, err := s.GetTraces(ctx, task.ID)
|
||
if err != nil {
|
||
t.Fatalf("get traces: %v", err)
|
||
}
|
||
var reviewTraces int
|
||
for _, tr := range traces {
|
||
if tr.Agent == "reviewer" {
|
||
reviewTraces++
|
||
}
|
||
}
|
||
if reviewTraces != 1 {
|
||
t.Errorf("reviewer запускался %d раз, want 1 (не должно быть retry второй сессии)", reviewTraces)
|
||
}
|
||
}
|
||
|
||
func TestWorkerTimeout(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
task := createReadyTask(t, s, "slow")
|
||
|
||
runner := &mockRunnerWorker{result: &opencode.Result{RC: -1, Stdout: ""}}
|
||
w := &Worker{
|
||
Store: s,
|
||
Runner: runner,
|
||
Worktree: t.TempDir(),
|
||
}
|
||
seedFakeRepo(t, w.Worktree, "slow")
|
||
|
||
ctx := context.Background()
|
||
_ = w.runTask(ctx, task)
|
||
|
||
task, err := s.GetTask(ctx, task.ID)
|
||
if err != nil {
|
||
t.Fatalf("get task: %v", err)
|
||
}
|
||
if task.Status != storage.StatusTimeout {
|
||
t.Errorf("status = %q, want timeout", task.Status)
|
||
}
|
||
|
||
traces, err := s.GetTraces(ctx, task.ID)
|
||
if err != nil {
|
||
t.Fatalf("get traces: %v", err)
|
||
}
|
||
if len(traces) != 2 {
|
||
t.Fatalf("got %d traces, want 2 (dev + postmortem)", len(traces))
|
||
}
|
||
if traces[0].Status != storage.TraceTimeout {
|
||
t.Errorf("trace[0] status = %q, want timeout", traces[0].Status)
|
||
}
|
||
if traces[1].Agent != postMortemAgent {
|
||
t.Errorf("trace[1] agent = %q, want postmortem", traces[1].Agent)
|
||
}
|
||
if traces[1].Status != storage.TraceSuccess {
|
||
t.Errorf("trace[1] status = %q, want success", traces[1].Status)
|
||
}
|
||
if runner.postMortemCount != 1 {
|
||
t.Errorf("postmortem запускался %d раз, want 1", runner.postMortemCount)
|
||
}
|
||
}
|
||
|
||
func TestWorkerSpawnError(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
task := createReadyTask(t, s, "spawn-fail")
|
||
|
||
runner := &mockRunnerWorker{err: errors.New("opencode not found")}
|
||
w := &Worker{
|
||
Store: s,
|
||
Runner: runner,
|
||
Worktree: t.TempDir(),
|
||
}
|
||
seedFakeRepo(t, w.Worktree, "spawn-fail")
|
||
|
||
ctx := context.Background()
|
||
_ = w.runTask(ctx, task)
|
||
|
||
task, err := s.GetTask(ctx, task.ID)
|
||
if err != nil {
|
||
t.Fatalf("get task: %v", err)
|
||
}
|
||
if task.Status != storage.StatusFailed {
|
||
t.Errorf("status = %q, want failed", task.Status)
|
||
}
|
||
|
||
traces, err := s.GetTraces(ctx, task.ID)
|
||
if err != nil {
|
||
t.Fatalf("get traces: %v", err)
|
||
}
|
||
if len(traces) != 2 {
|
||
t.Fatalf("got %d traces, want 2 (dev-failed + postmortem)", len(traces))
|
||
}
|
||
if traces[0].Status != storage.TraceFailed {
|
||
t.Errorf("trace[0] status = %q, want failed", traces[0].Status)
|
||
}
|
||
if traces[1].Agent != postMortemAgent {
|
||
t.Errorf("trace[1] agent = %q, want postmortem", traces[1].Agent)
|
||
}
|
||
if runner.postMortemCount != 1 {
|
||
t.Errorf("postmortem запускался %d раз, want 1", runner.postMortemCount)
|
||
}
|
||
}
|
||
|
||
func TestWorkerNonZeroExit(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
task := createReadyTask(t, s, "fail")
|
||
|
||
w := &Worker{
|
||
Store: s,
|
||
Runner: &mockRunnerWorker{result: &opencode.Result{RC: 7, Stdout: "error"}},
|
||
Worktree: t.TempDir(),
|
||
}
|
||
seedFakeRepo(t, w.Worktree, "fail")
|
||
|
||
ctx := context.Background()
|
||
_ = w.runTask(ctx, task)
|
||
|
||
task, err := s.GetTask(ctx, task.ID)
|
||
if err != nil {
|
||
t.Fatalf("get task: %v", err)
|
||
}
|
||
if task.Status != storage.StatusFailed {
|
||
t.Errorf("status = %q, want failed", task.Status)
|
||
}
|
||
|
||
traces, err := s.GetTraces(ctx, task.ID)
|
||
if err != nil {
|
||
t.Fatalf("get traces: %v", err)
|
||
}
|
||
if len(traces) != 2 {
|
||
t.Fatalf("got %d traces, want 2 (dev + postmortem)", len(traces))
|
||
}
|
||
if traces[0].Status != storage.TraceFailed {
|
||
t.Errorf("trace status = %q, want failed", traces[0].Status)
|
||
}
|
||
if traces[0].Output != "error" {
|
||
t.Errorf("output = %q, want error", traces[0].Output)
|
||
}
|
||
if traces[1].Agent != postMortemAgent {
|
||
t.Errorf("trace[1] agent = %q, want postmortem", traces[1].Agent)
|
||
}
|
||
}
|
||
|
||
func TestWorkerBadStatus(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
ctx := context.Background()
|
||
|
||
task := &storage.Task{ChatID: "tg://bad", Title: "bad-status", TaskTag: "bad"}
|
||
id, err := s.CreateTask(ctx, task)
|
||
if err != nil {
|
||
t.Fatalf("create task: %v", err)
|
||
}
|
||
task, _ = s.GetTask(ctx, id)
|
||
|
||
err = (&Worker{Store: s}).runTask(ctx, task)
|
||
if err == nil {
|
||
t.Fatal("expected error for non-ready task")
|
||
}
|
||
if !errors.Is(err, ErrLaunch) {
|
||
t.Errorf("err = %v, want W2", err)
|
||
}
|
||
}
|
||
|
||
func TestWorkerPromptRendered(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
task := createReadyTask(t, s, "prompt-test")
|
||
|
||
w := &Worker{
|
||
Store: s,
|
||
Runner: &mockRunnerWorker{result: &opencode.Result{RC: 0, Stdout: "ok", SessionID: "s"}},
|
||
Worktree: t.TempDir(),
|
||
}
|
||
seedFakeRepo(t, w.Worktree, "prompt-test")
|
||
|
||
ctx := context.Background()
|
||
_ = w.runTask(ctx, task)
|
||
|
||
traces, err := s.GetTraces(ctx, task.ID)
|
||
if err != nil {
|
||
t.Fatalf("get traces: %v", err)
|
||
}
|
||
if len(traces) == 0 {
|
||
t.Fatal("no traces")
|
||
}
|
||
tr := traces[0]
|
||
if tr.Prompt == "" {
|
||
t.Error("prompt empty — шаблон не срендерился")
|
||
}
|
||
if !strings.Contains(tr.Prompt, "prompt-test") {
|
||
t.Error("prompt не содержит название задачи")
|
||
}
|
||
if !strings.Contains(tr.Prompt, " - prompt-test") {
|
||
t.Error("prompt не содержит репозиторий prompt-test")
|
||
}
|
||
}
|
||
|
||
func TestFeatureBranchName(t *testing.T) {
|
||
cases := []struct {
|
||
tag string
|
||
want string
|
||
}{
|
||
{"", "feat/main"}, // пустой tag (старая задача) — не "feat/"
|
||
{" ", "feat/main"}, // только пробелы
|
||
{"3f2a1b", "feat/3f2a1b"}, // нормальный UUID
|
||
{" abc ", "feat/abc"}, // тримится
|
||
}
|
||
for _, c := range cases {
|
||
if got := featureBranchName(c.tag); got != c.want {
|
||
t.Errorf("featureBranchName(%q) = %q, want %q", c.tag, got, c.want)
|
||
}
|
||
}
|
||
}
|
||
|
||
func TestValidateRepoName(t *testing.T) {
|
||
valid := []string{"calc", "proj-a", "my.repo", "node_2", "kamelion/ratatoskr-go", "a/b/c"}
|
||
for _, r := range valid {
|
||
if err := validateRepoName(r); err != nil {
|
||
t.Errorf("validateRepoName(%q) = %v, want nil", r, err)
|
||
}
|
||
}
|
||
|
||
invalid := []string{"", "../etc", "etc/..", "a/../b", "..", ".", "/etc", "a//b", "./x"}
|
||
for _, r := range invalid {
|
||
if err := validateRepoName(r); err == nil {
|
||
t.Errorf("validateRepoName(%q) = nil, want E3", r)
|
||
} else if !errors.Is(err, ErrRepoPathHint) {
|
||
t.Errorf("validateRepoName(%q) err = %v, want E3", r, err)
|
||
}
|
||
}
|
||
}
|
||
|
||
func TestBuildCloneURL(t *testing.T) {
|
||
if got := buildCloneURL("http://gitea.hal9000.home", "proj-a"); got != "http://gitea.hal9000.home/proj-a.git" {
|
||
t.Errorf("buildCloneURL = %q", got)
|
||
}
|
||
if got := buildCloneURL("http://gitea.hal9000.home/", "proj-b"); got != "http://gitea.hal9000.home/proj-b.git" {
|
||
t.Errorf("buildCloneURL trailing slash = %q", got)
|
||
}
|
||
}
|
||
|
||
func TestPrepareRepos(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
_ = s
|
||
wt := t.TempDir()
|
||
w := &Worker{Worktree: wt, GitBaseURL: "http://gitea.hal9000.home"}
|
||
|
||
// seedFakeRepo уже создал .git — prepareRepos должен пройти без клона.
|
||
seedFakeRepo(t, wt, "proj-a")
|
||
if err := w.prepareRepos(context.Background(), []string{"proj-a"}); err != nil {
|
||
t.Fatalf("prepareRepos existing: %v", err)
|
||
}
|
||
|
||
// отсутствующий репо без git в PATH → E2 (clone упал), но не паника.
|
||
err := w.prepareRepos(context.Background(), []string{"missing"})
|
||
if err == nil {
|
||
t.Fatal("prepareRepos missing: expected error")
|
||
}
|
||
if !errors.Is(err, ErrClone) {
|
||
t.Errorf("prepareRepos missing err = %v, want E2", err)
|
||
}
|
||
}
|
||
|
||
// TestPrepareRepos_CreatesWorktreeDir проверяет, что prepareRepos создаёт сам
|
||
// каталог worktree на свежей машине (иначе git clone не может создать родителя,
|
||
// а decider падает с «не могу сменить папку на ./worktrees»).
|
||
func TestPrepareRepos_CreatesWorktreeDir(t *testing.T) {
|
||
// worktree ещё НЕ существует (родитель TempDir есть, сам каталог — нет).
|
||
base := t.TempDir()
|
||
wt := filepath.Join(base, "nested", "worktrees")
|
||
if _, err := os.Stat(wt); !os.IsNotExist(err) {
|
||
t.Fatalf("предусловие: %s должен отсутствовать, %v", wt, err)
|
||
}
|
||
|
||
w := &Worker{Worktree: wt, GitBaseURL: "http://gitea.hal9000.home"}
|
||
// Пустой список репо → prepareRepos всё равно создаёт каталог worktree.
|
||
if err := w.prepareRepos(context.Background(), nil); err != nil {
|
||
t.Fatalf("prepareRepos создал каталог: %v", err)
|
||
}
|
||
st, err := os.Stat(wt)
|
||
if err != nil {
|
||
t.Fatalf("worktree должен существовать после prepareRepos: %v", err)
|
||
}
|
||
if !st.IsDir() {
|
||
t.Fatalf("worktree не директория: %v", st.Mode())
|
||
}
|
||
}
|
||
|
||
func TestPrepareReposNonGitDir(t *testing.T) {
|
||
wt := t.TempDir()
|
||
w := &Worker{Worktree: wt}
|
||
|
||
// папка есть, но без .git → E4.
|
||
if err := os.MkdirAll(filepath.Join(wt, "plain"), 0o755); err != nil {
|
||
t.Fatal(err)
|
||
}
|
||
err := w.prepareRepos(context.Background(), []string{"plain"})
|
||
if err == nil {
|
||
t.Fatal("expected E4 error")
|
||
}
|
||
if !errors.Is(err, ErrRepoNotGit) {
|
||
t.Errorf("err = %v, want E4", err)
|
||
}
|
||
}
|
||
|
||
func TestWorkerStartStop(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
w := &Worker{
|
||
Store: s,
|
||
Runner: &mockRunnerWorker{result: &opencode.Result{RC: 0}},
|
||
Worktree: t.TempDir(),
|
||
MaxJobs: 1,
|
||
Interval: 50 * time.Millisecond,
|
||
}
|
||
|
||
ctx := context.Background()
|
||
w.Start(ctx)
|
||
|
||
time.Sleep(150 * time.Millisecond)
|
||
w.Stop()
|
||
}
|
||
|
||
// waitTaskStatus ждёт, пока задача достигнет статуса want. На Windows git-операции
|
||
// воркера заметно медленнее, чем на Linux, поэтому проверки в тестах не могут
|
||
// полагаться на фиксированные sleep'ы — только на polling до целевого статуса.
|
||
func waitTaskStatus(t *testing.T, ctx context.Context, s *storage.Storage, id int64, want storage.Status) {
|
||
t.Helper()
|
||
deadline := time.Now().Add(30 * time.Second)
|
||
for {
|
||
task, err := s.GetTask(ctx, id)
|
||
if err != nil {
|
||
t.Fatalf("get task %d: %v", id, err)
|
||
}
|
||
if task.Status == want {
|
||
return
|
||
}
|
||
switch task.Status {
|
||
case storage.StatusFailed, storage.StatusTimeout:
|
||
t.Fatalf("task %d: status %q, want %q", id, task.Status, want)
|
||
}
|
||
if time.Now().After(deadline) {
|
||
t.Fatalf("task %d: таймаут ожидания %q, последний статус %q", id, want, task.Status)
|
||
}
|
||
select {
|
||
case <-ctx.Done():
|
||
t.Fatalf("task %d: ctx done: %v", id, ctx.Err())
|
||
case <-time.After(50 * time.Millisecond):
|
||
}
|
||
}
|
||
}
|
||
|
||
func TestWorkerSemaphore(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
|
||
// создаём 2 ready-задачи
|
||
task1 := createReadyTask(t, s, "task-0")
|
||
task2 := createReadyTask(t, s, "task-1")
|
||
|
||
w := &Worker{
|
||
Store: s,
|
||
Runner: &mockRunnerWorker{result: &opencode.Result{RC: 0, Stdout: "ok"}},
|
||
Worktree: t.TempDir(),
|
||
MaxJobs: 1,
|
||
Interval: 50 * time.Millisecond,
|
||
}
|
||
seedFakeRepo(t, w.Worktree, "task-0")
|
||
seedFakeRepo(t, w.Worktree, "task-1")
|
||
w.sem = make(chan struct{}, 1)
|
||
w.sem <- struct{}{}
|
||
|
||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||
defer cancel()
|
||
|
||
// первый poll — запустит 1 задачу (макс. 1)
|
||
w.pollAndDispatch(ctx)
|
||
waitTaskStatus(t, ctx, s, task1.ID, storage.StatusSuccess)
|
||
|
||
// 1 должна быть success, 1 — всё ещё approved
|
||
success, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess})
|
||
approved, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusApproved})
|
||
if len(success) != 1 {
|
||
t.Errorf("success = %d, want 1 (approved=%d)", len(success), len(approved))
|
||
}
|
||
if len(approved) != 1 {
|
||
t.Errorf("approved = %d, want 1", len(approved))
|
||
}
|
||
|
||
// первая завершилась и вернула токен в сем — можем диспатчить вторую
|
||
w.pollAndDispatch(ctx)
|
||
waitTaskStatus(t, ctx, s, task2.ID, storage.StatusSuccess)
|
||
|
||
success, _ = s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess})
|
||
if len(success) != 2 {
|
||
t.Errorf("после освобождения слота success = %d, want 2", len(success))
|
||
}
|
||
}
|
||
|
||
// TestWorkerPostMortemSkippedOnSuccess — успешный трейд НЕ запускает постмортем:
|
||
// причина анализа — только failed/timeout.
|
||
func TestWorkerPostMortemSkippedOnSuccess(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
task := createReadyTask(t, s, "pm-ok")
|
||
|
||
runner := &mockRunnerWorker{result: &opencode.Result{RC: 0, Stdout: "done", SessionID: "s"}}
|
||
w := &Worker{
|
||
Store: s,
|
||
Runner: runner,
|
||
Worktree: t.TempDir(),
|
||
}
|
||
seedFakeRepo(t, w.Worktree, "pm-ok")
|
||
|
||
ctx := context.Background()
|
||
if err := w.runTask(ctx, task); err != nil {
|
||
t.Fatalf("runTask: %v", err)
|
||
}
|
||
task, _ = s.GetTask(ctx, task.ID)
|
||
if task.Status != storage.StatusSuccess {
|
||
t.Fatalf("status = %q, want success", task.Status)
|
||
}
|
||
if runner.postMortemCount != 0 {
|
||
t.Errorf("postmortem запускался %d раз, want 0 при success", runner.postMortemCount)
|
||
}
|
||
traces, _ := s.GetTraces(ctx, task.ID)
|
||
for _, tr := range traces {
|
||
if tr.Agent == postMortemAgent {
|
||
t.Errorf("есть неожиданный postmortem-trace при success")
|
||
}
|
||
}
|
||
}
|
||
|
||
// TestWorkerPostMortemFailureDoesNotChangeStatus — сбой самого постмортема не
|
||
// влияет на статус задачи (остаётся failed) и фиксируется как failed-трасса.
|
||
func TestWorkerPostMortemFailureDoesNotChangeStatus(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
task := createReadyTask(t, s, "pm-fail")
|
||
|
||
runner := &mockRunnerWorker{
|
||
// dev падает при спавне → failed; постмортем тоже падает.
|
||
err: errors.New("opencode not found"),
|
||
postMortemResult: &opencode.Result{RC: 1, Stdout: ""},
|
||
}
|
||
w := &Worker{
|
||
Store: s,
|
||
Runner: runner,
|
||
Worktree: t.TempDir(),
|
||
}
|
||
seedFakeRepo(t, w.Worktree, "pm-fail")
|
||
|
||
ctx := context.Background()
|
||
_ = w.runTask(ctx, task)
|
||
|
||
task, _ = s.GetTask(ctx, task.ID)
|
||
if task.Status != storage.StatusFailed {
|
||
t.Errorf("status = %q, want failed (постмортем не должен менять статус)", task.Status)
|
||
}
|
||
traces, err := s.GetTraces(ctx, task.ID)
|
||
if err != nil {
|
||
t.Fatalf("get traces: %v", err)
|
||
}
|
||
if len(traces) != 2 {
|
||
t.Fatalf("traces = %d, want 2 (dev-failed + postmortem)", len(traces))
|
||
}
|
||
if traces[1].Agent != postMortemAgent {
|
||
t.Errorf("trace[1] agent = %q, want postmortem", traces[1].Agent)
|
||
}
|
||
if traces[1].Status != storage.TraceFailed {
|
||
t.Errorf("trace[1] status = %q, want failed (сбой постмортема)", traces[1].Status)
|
||
}
|
||
}
|
||
|
||
// TestWorkerPostMortemNotRepeated — если у задачи уже есть postmortem-trace
|
||
// (например, от прошлого прогона), повторный постмортем не запускается.
|
||
func TestWorkerPostMortemNotRepeated(t *testing.T) {
|
||
s := setupWorkerDB(t)
|
||
task := createReadyTask(t, s, "pm-repeat")
|
||
|
||
ctx := context.Background()
|
||
if _, err := s.AppendTrace(ctx, &storage.Trace{
|
||
TaskID: task.ID,
|
||
Agent: postMortemAgent,
|
||
Prompt: "старый анализ",
|
||
Output: "старое резюме",
|
||
}); err != nil {
|
||
t.Fatalf("seed postmortem trace: %v", err)
|
||
}
|
||
|
||
runner := &mockRunnerWorker{err: errors.New("opencode not found")}
|
||
w := &Worker{
|
||
Store: s,
|
||
Runner: runner,
|
||
Worktree: t.TempDir(),
|
||
}
|
||
seedFakeRepo(t, w.Worktree, "pm-repeat")
|
||
|
||
_ = w.runTask(ctx, task)
|
||
|
||
task, _ = s.GetTask(ctx, task.ID)
|
||
if task.Status != storage.StatusFailed {
|
||
t.Fatalf("status = %q, want failed", task.Status)
|
||
}
|
||
if runner.postMortemCount != 0 {
|
||
t.Errorf("postmortem запускался %d раз, want 0 (уже был trace)", runner.postMortemCount)
|
||
}
|
||
// количество postmortem-трасс не выросло
|
||
traces, _ := s.GetTraces(ctx, task.ID)
|
||
var pm int
|
||
for _, tr := range traces {
|
||
if tr.Agent == postMortemAgent {
|
||
pm++
|
||
}
|
||
}
|
||
if pm != 1 {
|
||
t.Errorf("postmortem-трасс = %d, want 1 (без дубля)", pm)
|
||
}
|
||
}
|
||
|
||
// TestRenderPostMortemPrompt — промпт постмортема включает задачу и сессии.
|
||
func TestRenderPostMortemPrompt(t *testing.T) {
|
||
tr := storage.Trace{
|
||
Agent: "dev",
|
||
Status: storage.TraceTimeout,
|
||
Prompt: "промпт dev",
|
||
Output: "вывод dev",
|
||
}
|
||
prompt, err := RenderPostMortemPrompt(PostMortemPromptData{
|
||
Title: "Таймаут-задача",
|
||
Goal: "сделать",
|
||
Repos: []string{"calc"},
|
||
AC: "работает",
|
||
Status: storage.StatusTimeout,
|
||
Sessions: postMortemsText([]storage.Trace{tr}),
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("render: %v", err)
|
||
}
|
||
for _, want := range []string{"Таймаут-задача", "timeout", "=== Агент: dev", "промпт dev", "вывод dev"} {
|
||
if !strings.Contains(prompt, want) {
|
||
t.Errorf("промпт не содержит %q", want)
|
||
}
|
||
}
|
||
}
|