package worker import ( "context" "fmt" "log" "os" "os/exec" "path/filepath" "strings" "time" "github.com/kamelion/ratatoskr-go/internal/events" "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 // Notifier — механизм отправки авто-уведомлений владельцу задачи во время // выполнения. В проде реализуется *app.App через chat.Router.Send (см. // internal/app/app.go → App.Notify); в тестах worker подменяется фейковым // нотифаером. nil — уведомления выключены (ничего не отправляется). type Notifier interface { Notify(ctx context.Context, taskID int64, chatID, text string) 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 // Notify — нотифаер авто-уведомлений владельцу задачи (статусы + хендоффы // dev↔reviewer). nil — уведомления выключены. Notify Notifier sem chan struct{} // семафор cancel context.CancelFunc // подменяемый poll для тестов pollFn PollTaskFunc // Events — издатель доменных событий для UI. nil — события выключены. Events events.Publisher } // publish отправляет доменное событие, если задан издатель. func (w *Worker) publish(e events.Event) { if w.Events != nil { w.Events.Publish(e) } } // setStatus переводит задачу в новый статус: сохраняет в БД и публикует событие. func (w *Worker) setStatus(ctx context.Context, task *storage.Task, to storage.Status) error { from := task.Status task.Status = to if err := w.Store.UpdateTask(ctx, task); err != nil { return err } w.publish(events.TaskStatusChanged{ID: task.ID, From: from, To: to}) return nil } // 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) } // notify отправляет авто-уведомление владельцу задачи, если нотифаер задан. func (w *Worker) notify(ctx context.Context, task *storage.Task, text string) { if w.Notify == nil { return } if err := w.Notify.Notify(ctx, task.ID, task.ChatID, text); err != nil { log.Printf("worker: task %d: уведомление: %v", task.ID, err) } } // notifyStatus — уведомление о смене статуса задачи (номер задачи + статус). func (w *Worker) notifyStatus(ctx context.Context, task *storage.Task, s storage.Status) { w.notify(ctx, task, fmt.Sprintf("Задача #%d: %s", task.ID, s)) } // notifyHandoff — уведомление о передаче задачи между агентами конвейера // на заданной итерации (1-based). func (w *Worker) notifyHandoff(ctx context.Context, task *storage.Task, from, to string, iteration int) { w.notify(ctx, task, fmt.Sprintf("Задача #%d: %s → %s (итерация %d)", task.ID, from, to, iteration)) } // 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 (после валидации — чтобы плохие имена не жгли состояние) if err := w.setStatus(ctx, task, storage.StatusRunning); err != nil { return fmt.Errorf("%w: set running: %v", ErrUpdate, err) } w.notifyStatus(ctx, task, storage.StatusRunning) // 2b. клонируем недостающие репозитории в общий каталог. if err := w.prepareRepos(ctx, repos); err != nil { w.failTask(ctx, task) return fmt.Errorf("%w: %v", ErrClone, err) } // 2c. создаём feature-ветку от свежайшего origin/ в каждом репо. 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: if e := w.setStatus(ctx, task, storage.StatusTimeout); e != nil { return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e) } w.notifyStatus(ctx, task, storage.StatusTimeout) w.finalizeTrace(ctx, traceID, storage.TraceTimeout, output) return nil default: if e := w.setStatus(ctx, task, storage.StatusFailed); e != nil { return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e) } w.notifyStatus(ctx, task, storage.StatusFailed) w.finalizeTrace(ctx, traceID, storage.TraceFailed, output) return nil } // dev завершился RC=0 → сохраняем успех трассы dev. w.finalizeTrace(ctx, traceID, storage.TraceSuccess, output) // уведомляем пользователя о передаче dev → reviewer на ревью. w.notifyHandoff(ctx, task, "dev", "reviewer", iter+1) // 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 с объяснением. explain := "reviewer вернул невалидный/пустой вердикт (даже после повтора)." if e := w.setStatus(ctx, task, storage.StatusFailed); e != nil { return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e) } w.notifyStatus(ctx, task, storage.StatusFailed) 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 } if e := w.setStatus(ctx, task, storage.StatusSuccess); e != nil { return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e) } w.notifyStatus(ctx, task, storage.StatusSuccess) return nil } // Не пройдено: если есть итерации — dev дорабатывает. if iter+1 < maxReviewIterations { w.notify(ctx, task, fmt.Sprintf("Задача #%d: reviewer → dev на доработку (итерация %d)", task.ID, iter+1)) feedback = verdict.Comments continue } // Лимит исчерпан → failed с объяснением. if e := w.setStatus(ctx, task, storage.StatusFailed); e != nil { return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e) } w.notify(ctx, task, fmt.Sprintf("Задача #%d: failed — ревью не пройдено за %d итераций", task.ID, maxReviewIterations)) 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) { if e := w.setStatus(ctx, task, storage.StatusFailed); e != nil { log.Printf("worker: task %d: set failed: %v", task.ID, e) } w.notifyStatus(ctx, task, storage.StatusFailed) } // 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/. 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 }