Files
ratatoskr-go/internal/worker/worker.go
Hermes e055d1937c
All checks were successful
CI / test (push) Successful in 44s
CI / build-and-package (amd64, linux) (push) Successful in 43s
CI / build-and-package (amd64, windows) (push) Successful in 42s
fix: создавать каталог worktree в prepareRepos (decider: не могу сменить папку на ./worktrees)
2026-08-16 22:04:16 +05:00

402 lines
14 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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
sem chan struct{} // семафор
cancel context.CancelFunc
// подменяемый poll для тестов
pollFn PollTaskFunc
}
// 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.StatusReady,
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.StatusReady {
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(ctx, 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)
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).
func validateRepoName(repo string) error {
if repo == "" {
return fmt.Errorf("%w: пустое имя", ErrRepoPathHint)
}
if strings.Contains(repo, "/") || strings.Contains(repo, "..") {
return fmt.Errorf("%w: %q", ErrRepoPathHint, repo)
}
return nil
}