feat: воркер берёт задачи только после одобрения (approved)
Исправляет баг: воркер захватывал задачу на выполнение по статусу ready ещё до «создавай» (consent был заглушкой). Теперь: - новый статус approved: «создавай» → ready→approved; - воркер (pollAndDispatch + runTask) берёт ТОЛЬКО approved, ready = черновик готов, ждёт одобрения; - правка/текст в approved запрещены (финальное одобрение); - e2e-тест TestE2EWorkerDoesNotTakeUnconfirmed: в ready воркер задачу не трогает, запускает только после «создавай» → success; - обновлены все затронутые тесты (models/core/worker) и retry-фикстуры.
This commit is contained in:
@@ -233,14 +233,14 @@ func TestE2EWholeAppFromTaskSetup(t *testing.T) {
|
|||||||
t.Errorf("repos пусто, want [calc]")
|
t.Errorf("repos пусто, want [calc]")
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- 3. Согласие «создавай» → задача подтверждена (остаётся ready) ---
|
// --- 3. Согласие «создавай» → задача одобрена (approved) ---
|
||||||
fake.deliver(uid, "создавай")
|
fake.deliver(uid, "создавай")
|
||||||
task, err = a.Store.GetTask(ctx, task.ID)
|
task, err = a.Store.GetTask(ctx, task.ID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("get task: %v", err)
|
t.Fatalf("get task: %v", err)
|
||||||
}
|
}
|
||||||
if task.Status != storage.StatusReady {
|
if task.Status != storage.StatusApproved {
|
||||||
t.Errorf("status после создавай = %q, want ready", task.Status)
|
t.Errorf("status после создавай = %q, want approved", task.Status)
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- 4. Исполнение: воркер (poll) dev → reviewer → push ---
|
// --- 4. Исполнение: воркер (poll) dev → reviewer → push ---
|
||||||
@@ -331,9 +331,9 @@ func e2eChainToState(t *testing.T, store *storage.Storage, id int64, target stor
|
|||||||
var chain []storage.Status
|
var chain []storage.Status
|
||||||
switch target {
|
switch target {
|
||||||
case storage.StatusFailed:
|
case storage.StatusFailed:
|
||||||
chain = []storage.Status{storage.StatusCollecting, storage.StatusReady, storage.StatusRunning, storage.StatusFailed}
|
chain = []storage.Status{storage.StatusCollecting, storage.StatusReady, storage.StatusApproved, storage.StatusRunning, storage.StatusFailed}
|
||||||
case storage.StatusTimeout:
|
case storage.StatusTimeout:
|
||||||
chain = []storage.Status{storage.StatusCollecting, storage.StatusReady, storage.StatusRunning, storage.StatusTimeout}
|
chain = []storage.Status{storage.StatusCollecting, storage.StatusReady, storage.StatusApproved, storage.StatusRunning, storage.StatusTimeout}
|
||||||
default:
|
default:
|
||||||
t.Fatalf("e2eChainToState: неподдерживаемый target %q", target)
|
t.Fatalf("e2eChainToState: неподдерживаемый target %q", target)
|
||||||
}
|
}
|
||||||
@@ -418,14 +418,14 @@ func e2eRetryCommon(t *testing.T, a *App, fake *e2eChannel, worktree, uid string
|
|||||||
t.Errorf("retry: repos пусто, want [calc]")
|
t.Errorf("retry: repos пусто, want [calc]")
|
||||||
}
|
}
|
||||||
|
|
||||||
// 3. «создавай» → остаётся ready
|
// 3. «создавай» → approved (одобрено, воркер заберёт)
|
||||||
fake.deliver(chat.UserID(uid), "создавай")
|
fake.deliver(chat.UserID(uid), "создавай")
|
||||||
tk, err = a.Store.GetTask(ctx, taskID)
|
tk, err = a.Store.GetTask(ctx, taskID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("e2eRetryCommon: get task after создавай: %v", err)
|
t.Fatalf("e2eRetryCommon: get task after создавай: %v", err)
|
||||||
}
|
}
|
||||||
if tk.Status != storage.StatusReady {
|
if tk.Status != storage.StatusApproved {
|
||||||
t.Errorf("retry: status после создавай = %q, want ready", tk.Status)
|
t.Errorf("retry: status после создавай = %q, want approved", tk.Status)
|
||||||
}
|
}
|
||||||
|
|
||||||
// 4. воркер → success
|
// 4. воркер → success
|
||||||
@@ -628,4 +628,95 @@ func e2eRetryRejected(t *testing.T, a *App, fake *e2eChannel, uid string, taskID
|
|||||||
t.Errorf("retry(neg): нет понятного ответа «нельзя перезапустить», отправлено: %d сообщений", len(fake.sent))
|
t.Errorf("retry(neg): нет понятного ответа «нельзя перезапустить», отправлено: %d сообщений", len(fake.sent))
|
||||||
}
|
}
|
||||||
t.Logf("RETRY-REJECTED OK: задача #%d осталась %q", taskID, after.Status)
|
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)
|
||||||
}
|
}
|
||||||
@@ -52,7 +52,7 @@ func (c *Core) ProcessTurn(ctx context.Context, taskID int64, text string) (Resu
|
|||||||
return Result{}, err
|
return Result{}, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// 2. согласие в фазе ready → создание
|
// 2. согласие в фазе ready → одобрение
|
||||||
if task.Status == storage.StatusReady {
|
if task.Status == storage.StatusReady {
|
||||||
if isConsent(text) {
|
if isConsent(text) {
|
||||||
return c.handleConsent(ctx, task)
|
return c.handleConsent(ctx, task)
|
||||||
@@ -61,6 +61,17 @@ func (c *Core) ProcessTurn(ctx context.Context, taskID int64, text string) (Resu
|
|||||||
return c.handleEdit(ctx, task, text)
|
return c.handleEdit(ctx, task, text)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 2b. approved — финальное одобрение, правка запрещена.
|
||||||
|
// Воркер уже взял/заберёт задачу; текст не меняет статус.
|
||||||
|
if task.Status == storage.StatusApproved {
|
||||||
|
return Result{
|
||||||
|
Reply: "Задача уже одобрена и передана на выполнение. Следите за статусом: /status " + itoa(task.ID),
|
||||||
|
Action: "send",
|
||||||
|
TaskID: task.ID,
|
||||||
|
Status: task.Status,
|
||||||
|
}, nil
|
||||||
|
}
|
||||||
|
|
||||||
// 3. обычный ход: накопление + аналитик
|
// 3. обычный ход: накопление + аналитик
|
||||||
return c.handleTurn(ctx, task, text)
|
return c.handleTurn(ctx, task, text)
|
||||||
}
|
}
|
||||||
@@ -271,14 +282,15 @@ func (c *Core) notFoundReply(ctx context.Context, id int64, err error) (Result,
|
|||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// handleConsent создаёт задачу (статус ready → ...). Пока — подтверждение готовности.
|
// handleConsent одобряет задачу: ready → approved (финальное одобрение,
|
||||||
|
// после которого воркер забирает задачу на выполнение).
|
||||||
func (c *Core) handleConsent(ctx context.Context, task *storage.Task) (Result, error) {
|
func (c *Core) handleConsent(ctx context.Context, task *storage.Task) (Result, error) {
|
||||||
task.Status = storage.StatusReady
|
task.Status = storage.StatusApproved
|
||||||
if err := c.Store.UpdateTask(ctx, task); err != nil {
|
if err := c.Store.UpdateTask(ctx, task); err != nil {
|
||||||
return Result{}, err
|
return Result{}, err
|
||||||
}
|
}
|
||||||
return Result{
|
return Result{
|
||||||
Reply: "✅ Задача #" + itoa(task.ID) + " готова к запуску.",
|
Reply: "✅ Задача #" + itoa(task.ID) + " одобрена. Запускаю выполнение.",
|
||||||
Action: "created:" + itoa(task.ID),
|
Action: "created:" + itoa(task.ID),
|
||||||
TaskID: task.ID,
|
TaskID: task.ID,
|
||||||
Status: task.Status,
|
Status: task.Status,
|
||||||
|
|||||||
@@ -160,6 +160,10 @@ func TestConsentInReady(t *testing.T) {
|
|||||||
if res.Action != "created:"+itoa(id) {
|
if res.Action != "created:"+itoa(id) {
|
||||||
t.Fatalf("action = %q, want created:%d", res.Action, id)
|
t.Fatalf("action = %q, want created:%d", res.Action, id)
|
||||||
}
|
}
|
||||||
|
task, _ := store.GetTask(ctx, id)
|
||||||
|
if task.Status != storage.StatusApproved {
|
||||||
|
t.Fatalf("status после создавай = %s, want approved", task.Status)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestEditInReadyGoesCollecting(t *testing.T) {
|
func TestEditInReadyGoesCollecting(t *testing.T) {
|
||||||
|
|||||||
@@ -8,7 +8,8 @@ type Status string
|
|||||||
const (
|
const (
|
||||||
StatusDraft Status = "draft" // только что создана
|
StatusDraft Status = "draft" // только что создана
|
||||||
StatusCollecting Status = "collecting" // аналитик собирает детали
|
StatusCollecting Status = "collecting" // аналитик собирает детали
|
||||||
StatusReady Status = "ready" // черновик готов, ждёт запуска
|
StatusReady Status = "ready" // черновик готов, ждёт одобрения пользователя
|
||||||
|
StatusApproved Status = "approved" // пользователь одобрил («создавай») — воркер берёт в работу
|
||||||
StatusRunning Status = "running" // opencode работает
|
StatusRunning Status = "running" // opencode работает
|
||||||
StatusSuccess Status = "success" // задача выполнена
|
StatusSuccess Status = "success" // задача выполнена
|
||||||
StatusFailed Status = "failed" // ошибка выполнения
|
StatusFailed Status = "failed" // ошибка выполнения
|
||||||
@@ -20,23 +21,26 @@ const (
|
|||||||
|
|
||||||
// AllStatuses — все возможные статусы для валидации.
|
// AllStatuses — все возможные статусы для валидации.
|
||||||
var AllStatuses = []Status{
|
var AllStatuses = []Status{
|
||||||
StatusDraft, StatusCollecting, StatusReady,
|
StatusDraft, StatusCollecting, StatusReady, StatusApproved,
|
||||||
StatusRunning, StatusSuccess, StatusFailed, StatusTimeout,
|
StatusRunning, StatusSuccess, StatusFailed, StatusTimeout,
|
||||||
StatusCancelled, StatusAborted, StatusClosed,
|
StatusCancelled, StatusAborted, StatusClosed,
|
||||||
}
|
}
|
||||||
|
|
||||||
// validTransitions задаёт разрешённые переходы статусов.
|
// validTransitions задаёт разрешённые переходы статусов.
|
||||||
var validTransitions = map[Status][]Status{
|
var validTransitions = map[Status][]Status{
|
||||||
StatusDraft: {StatusCollecting, StatusCancelled, StatusAborted},
|
StatusDraft: {StatusCollecting, StatusCancelled, StatusAborted},
|
||||||
StatusCollecting: {StatusReady, StatusDraft, StatusCancelled, StatusAborted},
|
StatusCollecting: {StatusReady, StatusDraft, StatusCancelled, StatusAborted},
|
||||||
StatusReady: {StatusRunning, StatusCancelled, StatusAborted, StatusClosed, StatusCollecting}, // правка готового
|
// ready — черновик готов: «создавай» → approved, либо правка/отмена/закрытие.
|
||||||
StatusRunning: {StatusSuccess, StatusFailed, StatusTimeout, StatusCancelled},
|
StatusReady: {StatusApproved, StatusCancelled, StatusAborted, StatusClosed, StatusCollecting},
|
||||||
StatusSuccess: {StatusClosed},
|
// approved — финальное одобрение: воркер берёт в running, либо отмена/сбой/закрытие.
|
||||||
StatusFailed: {StatusReady, StatusClosed, StatusCancelled, StatusCollecting}, // retry: перезапуск сбора
|
StatusApproved: {StatusRunning, StatusCancelled, StatusAborted, StatusClosed},
|
||||||
StatusTimeout: {StatusReady, StatusClosed, StatusCancelled, StatusCollecting}, // retry: перезапуск сбора
|
StatusRunning: {StatusSuccess, StatusFailed, StatusTimeout, StatusCancelled},
|
||||||
StatusCancelled: {StatusClosed},
|
StatusSuccess: {StatusClosed},
|
||||||
StatusAborted: {StatusClosed},
|
StatusFailed: {StatusReady, StatusClosed, StatusCancelled, StatusCollecting}, // retry: перезапуск сбора
|
||||||
StatusClosed: {}, // терминальный
|
StatusTimeout: {StatusReady, StatusClosed, StatusCancelled, StatusCollecting}, // retry: перезапуск сбора
|
||||||
|
StatusCancelled: {StatusClosed},
|
||||||
|
StatusAborted: {StatusClosed},
|
||||||
|
StatusClosed: {}, // терминальный
|
||||||
}
|
}
|
||||||
|
|
||||||
// IsValidTransition проверяет, допустим ли переход from → to.
|
// IsValidTransition проверяет, допустим ли переход from → to.
|
||||||
|
|||||||
@@ -25,6 +25,12 @@ func TestIsValidTransition(t *testing.T) {
|
|||||||
{StatusTimeout, StatusReady, true}, // retry
|
{StatusTimeout, StatusReady, true}, // retry
|
||||||
{StatusTimeout, StatusCollecting, true}, // retry: перезапуск сбора
|
{StatusTimeout, StatusCollecting, true}, // retry: перезапуск сбора
|
||||||
{StatusTimeout, StatusRunning, false},
|
{StatusTimeout, StatusRunning, false},
|
||||||
|
{StatusReady, StatusApproved, true}, // «создавай» → одобрено
|
||||||
|
{StatusReady, StatusRunning, false}, // без одобрения воркер не запускает
|
||||||
|
{StatusApproved, StatusRunning, true}, // воркер берёт approved
|
||||||
|
{StatusApproved, StatusCancelled, true}, // отмена до запуска
|
||||||
|
{StatusApproved, StatusReady, false}, // финал: назад нельзя
|
||||||
|
{StatusApproved, StatusCollecting, false},
|
||||||
{StatusClosed, StatusDraft, false},
|
{StatusClosed, StatusDraft, false},
|
||||||
{StatusClosed, StatusRunning, false},
|
{StatusClosed, StatusRunning, false},
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -113,7 +113,7 @@ func (w *Worker) pollAndDispatch(ctx context.Context) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
tasks, err := w.Store.ListTasks(ctx, storage.TaskFilter{
|
tasks, err := w.Store.ListTasks(ctx, storage.TaskFilter{
|
||||||
Status: storage.StatusReady,
|
Status: storage.StatusApproved,
|
||||||
Limit: slots,
|
Limit: slots,
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -142,7 +142,7 @@ func (w *Worker) pollAndDispatch(ctx context.Context) error {
|
|||||||
// dev дорабатывает по комментариям; прошло → push ветки + success.
|
// dev дорабатывает по комментариям; прошло → push ветки + success.
|
||||||
func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
||||||
// 1. проверяем статус
|
// 1. проверяем статус
|
||||||
if task.Status != storage.StatusReady {
|
if task.Status != storage.StatusApproved {
|
||||||
return fmt.Errorf("%w: task %d status=%q", ErrLaunch, task.ID, task.Status)
|
return fmt.Errorf("%w: task %d status=%q", ErrLaunch, task.ID, task.Status)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -95,6 +95,10 @@ func createReadyTask(t *testing.T, s *storage.Storage, title string) *storage.Ta
|
|||||||
if err := s.UpdateTask(ctx, task); err != nil {
|
if err := s.UpdateTask(ctx, task); err != nil {
|
||||||
t.Fatalf("set ready: %v", err)
|
t.Fatalf("set ready: %v", err)
|
||||||
}
|
}
|
||||||
|
task.Status = storage.StatusApproved
|
||||||
|
if err := s.UpdateTask(ctx, task); err != nil {
|
||||||
|
t.Fatalf("set approved: %v", err)
|
||||||
|
}
|
||||||
task, _ = s.GetTask(ctx, id)
|
task, _ = s.GetTask(ctx, id)
|
||||||
return task
|
return task
|
||||||
}
|
}
|
||||||
@@ -625,14 +629,14 @@ func TestWorkerSemaphore(t *testing.T) {
|
|||||||
w.pollAndDispatch(ctx)
|
w.pollAndDispatch(ctx)
|
||||||
time.Sleep(200 * time.Millisecond)
|
time.Sleep(200 * time.Millisecond)
|
||||||
|
|
||||||
// 1 должна быть success, 1 — всё ещё ready
|
// 1 должна быть success, 1 — всё ещё approved
|
||||||
success, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess})
|
success, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess})
|
||||||
ready, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusReady})
|
approved, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusApproved})
|
||||||
if len(success) != 1 {
|
if len(success) != 1 {
|
||||||
t.Errorf("success = %d, want 1 (ready=%d)", len(success), len(ready))
|
t.Errorf("success = %d, want 1 (approved=%d)", len(success), len(approved))
|
||||||
}
|
}
|
||||||
if len(ready) != 1 {
|
if len(approved) != 1 {
|
||||||
t.Errorf("ready = %d, want 1", len(ready))
|
t.Errorf("approved = %d, want 1", len(approved))
|
||||||
}
|
}
|
||||||
|
|
||||||
// первая завершилась и вернула токен в сем — можем диспатчить вторую
|
// первая завершилась и вернула токен в сем — можем диспатчить вторую
|
||||||
|
|||||||
Reference in New Issue
Block a user