Исправляет баг: воркер захватывал задачу на выполнение по статусу ready ещё до «создавай» (consent был заглушкой). Теперь: - новый статус approved: «создавай» → ready→approved; - воркер (pollAndDispatch + runTask) берёт ТОЛЬКО approved, ready = черновик готов, ждёт одобрения; - правка/текст в approved запрещены (финальное одобрение); - e2e-тест TestE2EWorkerDoesNotTakeUnconfirmed: в ready воркер задачу не трогает, запускает только после «создавай» → success; - обновлены все затронутые тесты (models/core/worker) и retry-фикстуры.
433 lines
15 KiB
Go
433 lines
15 KiB
Go
package worker
|
||
|
||
import (
|
||
"context"
|
||
"fmt"
|
||
"log"
|
||
"os"
|
||
"os/exec"
|
||
"path/filepath"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/kamelion/ratatoskr-go/internal/opencode"
|
||
"github.com/kamelion/ratatoskr-go/internal/storage"
|
||
)
|
||
|
||
// OpenCodeRunner — интерфейс для opencode (подменяемый в тестах).
|
||
type OpenCodeRunner interface {
|
||
Run(ctx context.Context, prompt, cwd, agent, sessionID string) (*opencode.Result, error)
|
||
}
|
||
|
||
// PollTaskFunc — callback для обработки готовой задачи (подменяемый в тестах).
|
||
type PollTaskFunc func(ctx context.Context) error
|
||
|
||
// Worker — планировщик, запускающий готовые задачи (status=ready → running → success/failed/timeout).
|
||
type Worker struct {
|
||
Store *storage.Storage
|
||
Runner OpenCodeRunner
|
||
Worktree string // общий каталог, репозитории вкладываются в него по имени
|
||
Agent string // default "dev"
|
||
Interval time.Duration // интервал опроса БД
|
||
MaxJobs int // макс. параллельных задач
|
||
|
||
// Git — источник репозиториев для клонирования.
|
||
GitBaseURL string
|
||
GitToken string
|
||
|
||
// Live — опциональный журнал живых сессий (наблюдение /status N).
|
||
// Через него Runner пишет live-шаги задачи; nil — наблюдение выключено.
|
||
Live *opencode.LiveRegistry
|
||
|
||
sem chan struct{} // семафор
|
||
cancel context.CancelFunc
|
||
|
||
// подменяемый poll для тестов
|
||
pollFn PollTaskFunc
|
||
}
|
||
|
||
// runCtx оборачивает контекст запуска субагента, привязывая живое наблюдение
|
||
// сессии задачи (если Live-журнал включён). Возвращает ctx без изменений при nil Live.
|
||
func (w *Worker) runCtx(ctx context.Context, taskID int64) context.Context {
|
||
if w.Live == nil {
|
||
return ctx
|
||
}
|
||
return opencode.WithLive(ctx, w.Live, taskID)
|
||
}
|
||
|
||
// Start запускает цикл опроса в фоновой горутине.
|
||
func (w *Worker) Start(ctx context.Context) {
|
||
if w.Agent == "" {
|
||
w.Agent = "dev"
|
||
}
|
||
if w.Interval <= 0 {
|
||
w.Interval = 5 * time.Second
|
||
}
|
||
if w.MaxJobs <= 0 {
|
||
w.MaxJobs = 2
|
||
}
|
||
w.sem = make(chan struct{}, w.MaxJobs)
|
||
// заполняем семафор токенами
|
||
for i := 0; i < w.MaxJobs; i++ {
|
||
w.sem <- struct{}{}
|
||
}
|
||
|
||
ctx, w.cancel = context.WithCancel(ctx)
|
||
|
||
pollFn := w.pollFn
|
||
if pollFn == nil {
|
||
pollFn = w.pollAndDispatch
|
||
}
|
||
|
||
go func() {
|
||
// первый poll сразу
|
||
_ = pollFn(ctx)
|
||
|
||
ticker := time.NewTicker(w.Interval)
|
||
defer ticker.Stop()
|
||
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
log.Print("worker: stopped")
|
||
return
|
||
case <-ticker.C:
|
||
_ = pollFn(ctx)
|
||
}
|
||
}
|
||
}()
|
||
}
|
||
|
||
// Stop останавливает воркер (отменяет контекст → убивает активные задачи).
|
||
func (w *Worker) Stop() {
|
||
if w.cancel != nil {
|
||
w.cancel()
|
||
}
|
||
}
|
||
|
||
// pollAndDispatch ищет готовые задачи, запускает их в пределах свободных слотов.
|
||
func (w *Worker) pollAndDispatch(ctx context.Context) error {
|
||
slots := len(w.sem)
|
||
if slots == 0 {
|
||
return nil
|
||
}
|
||
|
||
tasks, err := w.Store.ListTasks(ctx, storage.TaskFilter{
|
||
Status: storage.StatusApproved,
|
||
Limit: slots,
|
||
})
|
||
if err != nil {
|
||
return fmt.Errorf("%w: %v", ErrPoll, err)
|
||
}
|
||
|
||
for _, t := range tasks {
|
||
select {
|
||
case <-ctx.Done():
|
||
return ctx.Err()
|
||
case slot := <-w.sem:
|
||
task := t
|
||
go func() {
|
||
defer func() { w.sem <- slot }()
|
||
if err := w.runTask(ctx, task); err != nil {
|
||
log.Printf("worker: task %d: %v", task.ID, err)
|
||
}
|
||
}()
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// runTask выполняет одну задачу: готовит репозитории, создаёт feature-ветку,
|
||
// dev-агент реализует, reviewer строго проверяет весь дифф ветки; при не-проходе
|
||
// dev дорабатывает по комментариям; прошло → push ветки + success.
|
||
func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
||
// 1. проверяем статус
|
||
if task.Status != storage.StatusApproved {
|
||
return fmt.Errorf("%w: task %d status=%q", ErrLaunch, task.ID, task.Status)
|
||
}
|
||
|
||
repos := task.EffectiveRepos()
|
||
if len(repos) == 0 {
|
||
return fmt.Errorf("%w: task %d: %v", ErrLaunch, task.ID, ErrNoRepos)
|
||
}
|
||
|
||
// 1b. проверяем имена репо (E3): не допускаем путь-escape.
|
||
for _, r := range repos {
|
||
if err := validateRepoName(r); err != nil {
|
||
return fmt.Errorf("%w: task %d repo %q: %w", ErrLaunch, task.ID, r, err)
|
||
}
|
||
}
|
||
|
||
// 2. ставим running (после валидации — чтобы плохие имена не жгли состояние)
|
||
task.Status = storage.StatusRunning
|
||
if err := w.Store.UpdateTask(ctx, task); err != nil {
|
||
return fmt.Errorf("%w: set running: %v", ErrUpdate, err)
|
||
}
|
||
|
||
// 2b. клонируем недостающие репозитории в общий каталог.
|
||
if err := w.prepareRepos(ctx, repos); err != nil {
|
||
w.failTask(ctx, task)
|
||
return fmt.Errorf("%w: %v", ErrClone, err)
|
||
}
|
||
|
||
// 2c. создаём feature-ветку от свежайшего origin/<baseBranch> в каждом репо.
|
||
branch := featureBranchName(task.TaskTag)
|
||
for _, r := range repos {
|
||
if err := w.ensureBranch(ctx, w.repoDirOf(r), branch); err != nil {
|
||
w.failTask(ctx, task)
|
||
return fmt.Errorf("%w: create branch %s in %s: %v", ErrClone, branch, r, err)
|
||
}
|
||
}
|
||
|
||
// 3. cwd — общий каталог (вариант A: один dev видит все репозитории).
|
||
cwd := w.Worktree
|
||
|
||
// Цикл dev → review, до maxReviewIterations.
|
||
var feedback []string
|
||
|
||
for iter := 0; ; iter++ {
|
||
// 3. рендерим промпт dev (с feedback на повторных итерациях)
|
||
devData := DevPromptData{
|
||
Title: task.Title,
|
||
Goal: task.Goal,
|
||
Repos: repos,
|
||
Why: task.Why,
|
||
AC: task.AC,
|
||
Branch: branch,
|
||
ReviewFeedback: reviewFeedbackList(branch, feedback),
|
||
}
|
||
prompt, perr := RenderDevPrompt(devData)
|
||
if perr != nil {
|
||
return fmt.Errorf("%w: render prompt: %v", ErrTrace, perr)
|
||
}
|
||
|
||
// 4b. создаём трассу dev
|
||
trace := &storage.Trace{TaskID: task.ID, Agent: w.Agent, Prompt: prompt}
|
||
traceID, tErr := w.Store.AppendTrace(ctx, trace)
|
||
if tErr != nil {
|
||
return fmt.Errorf("%w: create: %v", ErrTrace, tErr)
|
||
}
|
||
|
||
// 5. запускаем dev-агента (fresh сессия в текущей ветке)
|
||
res, resErr := w.Runner.Run(w.runCtx(ctx, task.ID), prompt, cwd, w.Agent, "")
|
||
if resErr != nil {
|
||
// O1 ErrSpawn — не смог запустить бинарь
|
||
w.failTask(ctx, task)
|
||
w.finalizeTrace(ctx, traceID, storage.TraceFailed, resErr.Error())
|
||
return fmt.Errorf("%w: spawn: %v", ErrLaunch, resErr)
|
||
}
|
||
if res.SessionID != "" {
|
||
_ = w.Store.UpdateTraceSessionID(ctx, traceID, res.SessionID)
|
||
}
|
||
|
||
output := res.Stdout
|
||
|
||
// 5b. dev не завершился успешно (RC!=0) → фиксируем без ревью.
|
||
switch res.RC {
|
||
case 0:
|
||
// продолжаем на ревью
|
||
case -1:
|
||
task.Status = storage.StatusTimeout
|
||
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
||
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
||
}
|
||
w.finalizeTrace(ctx, traceID, storage.TraceTimeout, output)
|
||
return nil
|
||
default:
|
||
task.Status = storage.StatusFailed
|
||
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
||
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
||
}
|
||
w.finalizeTrace(ctx, traceID, storage.TraceFailed, output)
|
||
return nil
|
||
}
|
||
|
||
// dev завершился RC=0 → сохраняем успех трассы dev.
|
||
w.finalizeTrace(ctx, traceID, storage.TraceSuccess, output)
|
||
|
||
// 8. РЕВЬЮ: собираем diff всей ветки, запускаем reviewer.
|
||
diffText, dErr := w.branchDiffAll(ctx, repos, branch)
|
||
if dErr != nil {
|
||
w.failTask(ctx, task)
|
||
return fmt.Errorf("%w: %v", ErrReviewDiff, dErr)
|
||
}
|
||
reviewPrompt, rErr := RenderReviewPrompt(ReviewPromptData{
|
||
Branch: branch,
|
||
AC: task.AC,
|
||
Diff: diffText,
|
||
})
|
||
if rErr != nil {
|
||
w.failTask(ctx, task)
|
||
return fmt.Errorf("%w: render review prompt: %v", ErrReviewTrace, rErr)
|
||
}
|
||
|
||
// Запускаем reviewer, с одним retry на невалидный JSON/вывод.
|
||
verdict, reviewOutput, reviewTraceID, rvErr := w.reviewWithRetry(ctx, task.ID, cwd, reviewPrompt)
|
||
if rvErr != nil {
|
||
w.failTask(ctx, task)
|
||
return rvErr
|
||
}
|
||
|
||
if verdict == nil {
|
||
// невалидный JSON даже после retry → failed с объяснением.
|
||
task.Status = storage.StatusFailed
|
||
explain := "reviewer вернул невалидный/пустой вердикт (даже после повтора)."
|
||
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
||
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
||
}
|
||
w.finalizeTrace(ctx, reviewTraceID, storage.TraceFailed, reviewOutput+"\n"+explain)
|
||
return nil
|
||
}
|
||
|
||
// Пройдено → push и success.
|
||
if verdict.Passed {
|
||
if pErr := w.pushBranches(ctx, repos, branch); pErr != nil {
|
||
w.failTask(ctx, task)
|
||
return pErr
|
||
}
|
||
task.Status = storage.StatusSuccess
|
||
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
||
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// Не пройдено: если есть итерации — dev дорабатывает.
|
||
if iter+1 < maxReviewIterations {
|
||
feedback = verdict.Comments
|
||
continue
|
||
}
|
||
|
||
// Лимит исчерпан → failed с объяснением.
|
||
task.Status = storage.StatusFailed
|
||
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
||
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
|
||
}
|
||
explain := fmt.Sprintf("Ревью не пройдено за %d итераций.", maxReviewIterations)
|
||
final := reviewOutput + "\n" + explain
|
||
if e := w.Store.UpdateTraceOutput(ctx, reviewTraceID, final); e != nil {
|
||
log.Printf("worker: task %d: update review trace: %v", task.ID, e)
|
||
}
|
||
return nil
|
||
}
|
||
}
|
||
|
||
// reviewWithRetry запускает reviewer; при непарсируемом вердикте — один повтор.
|
||
// Возвращает (verdict, output, traceID). traceID — последней попытки ревью.
|
||
func (w *Worker) reviewWithRetry(ctx context.Context, taskID int64, cwd, prompt string) (*reviewVerdict, string, int64, error) {
|
||
v, out, tid, err := w.runReviewer(ctx, taskID, cwd, prompt)
|
||
if err != nil {
|
||
return nil, out, tid, err
|
||
}
|
||
if v != nil {
|
||
return v, out, tid, nil
|
||
}
|
||
// невалидный/пустой — один retry.
|
||
v2, out2, tid2, err2 := w.runReviewer(ctx, taskID, cwd, prompt)
|
||
if err2 != nil {
|
||
return nil, out2, tid2, err2
|
||
}
|
||
return v2, out2, tid2, nil
|
||
}
|
||
|
||
// failTask помечает задачу failed.
|
||
func (w *Worker) failTask(ctx context.Context, task *storage.Task) {
|
||
task.Status = storage.StatusFailed
|
||
if e := w.Store.UpdateTask(ctx, task); e != nil {
|
||
log.Printf("worker: task %d: set failed: %v", task.ID, e)
|
||
}
|
||
}
|
||
|
||
// finalizeTrace обновляет output и статус трассы.
|
||
func (w *Worker) finalizeTrace(ctx context.Context, traceID int64, status storage.TraceStatus, output string) {
|
||
if e := w.Store.UpdateTraceOutput(ctx, traceID, output); e != nil {
|
||
log.Printf("worker: update trace output %d: %v", traceID, e)
|
||
}
|
||
if e := w.Store.UpdateTraceStatus(ctx, traceID, status); e != nil {
|
||
log.Printf("worker: update trace status %d: %v", traceID, e)
|
||
}
|
||
}
|
||
|
||
// prepareRepos гарантирует наличие всех репозиториев в общем каталоге:
|
||
// папка есть → используем как есть; нет → git clone.
|
||
func (w *Worker) prepareRepos(ctx context.Context, repos []string) error {
|
||
// Общий каталог worktree может отсутствовать на свежей машине: git clone не
|
||
// создаёт родительский каталог, а decider стартует с cwd=Worktree. Создаём заранее.
|
||
if err := os.MkdirAll(w.Worktree, 0o755); err != nil {
|
||
return fmt.Errorf("%w: создать каталог worktree %s: %v", ErrClone, w.Worktree, err)
|
||
}
|
||
for _, r := range repos {
|
||
dst := filepath.Join(w.Worktree, r)
|
||
if _, err := os.Stat(dst); err == nil {
|
||
gitdir := filepath.Join(dst, ".git")
|
||
if _, e := os.Stat(gitdir); e != nil {
|
||
// E4: папка есть, но не git-репо — клонировать поверх нельзя.
|
||
return fmt.Errorf("%w: %s существует, но не git-репозиторий", ErrRepoNotGit, r)
|
||
}
|
||
continue // уже готово
|
||
}
|
||
|
||
if err := w.clone(ctx, r); err != nil {
|
||
return fmt.Errorf("%w: %v", ErrClone, err)
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// clone клонирует репозиторий r в ./worktrees/<r>.
|
||
func (w *Worker) clone(ctx context.Context, repo string) error {
|
||
base := strings.TrimRight(w.GitBaseURL, "/")
|
||
if base == "" {
|
||
return fmt.Errorf("git.base_url не задан в конфиге")
|
||
}
|
||
url := buildCloneURL(base, repo)
|
||
dst := filepath.Join(w.Worktree, repo)
|
||
|
||
// Для формата "owner/repo" git clone не создаёт промежуточный каталог {worktree}/owner
|
||
// до самого репозитория — поэтому создаём родителя явно. Worktree в целом гарантирован
|
||
// app.New, но вложенные сегменты owner здесь обязательны.
|
||
if parent := filepath.Dir(dst); parent != "." && parent != w.Worktree {
|
||
if err := os.MkdirAll(parent, 0o755); err != nil {
|
||
return fmt.Errorf("%w: создать %s: %v", ErrClone, parent, err)
|
||
}
|
||
}
|
||
|
||
args := []string{"clone"}
|
||
if w.GitToken != "" {
|
||
// https-базовый URL: кладём токен внутрь URL (для приватных репозиториев).
|
||
args = append(args, "--config", "http.extraHeader=Authorization: Bearer "+w.GitToken)
|
||
}
|
||
args = append(args, url, dst)
|
||
|
||
cmd := exec.CommandContext(ctx, "git", args...)
|
||
out, err := cmd.CombinedOutput()
|
||
if err != nil {
|
||
return fmt.Errorf("%w: git clone %s: %s", ErrClone, repo, strings.TrimSpace(string(out)))
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// buildCloneURL собирает URL клона из base_url и имени репозитория.
|
||
func buildCloneURL(base, repo string) string {
|
||
return strings.TrimRight(base, "/") + "/" + repo + ".git"
|
||
}
|
||
|
||
// validateRepoName отклоняет имена с пути-эскейпом (E3).
|
||
// Допускается формат "owner/repo" (один слэш) — клон ляжет в {worktree}/owner/repo,
|
||
// промежуточный каталог создаёт clone(). Путь с сегментами "." , ".." , пустыми или
|
||
// абсолютный — ошибка.
|
||
func validateRepoName(repo string) error {
|
||
if repo == "" {
|
||
return fmt.Errorf("%w: пустое имя", ErrRepoPathHint)
|
||
}
|
||
if filepath.IsAbs(repo) {
|
||
return fmt.Errorf("%w: %q", ErrRepoPathHint, repo)
|
||
}
|
||
for _, seg := range strings.Split(repo, "/") {
|
||
if seg == "" || seg == "." || seg == ".." {
|
||
return fmt.Errorf("%w: %q", ErrRepoPathHint, repo)
|
||
}
|
||
}
|
||
return nil
|
||
}
|