Files
ratatoskr-go/internal/worker/worker.go
ki.sagidullin 07bae203c3
Some checks failed
CI / test (pull_request) Failing after 35s
CI / build-and-package (amd64, linux) (pull_request) Successful in 40s
CI / build-and-package (amd64, windows) (pull_request) Successful in 43s
feat: авто-уведомления о статусах и хендоффах задачи
2026-08-17 23:32:26 +05:00

477 lines
18 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
// 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
}
// 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 (после валидации — чтобы плохие имена не жгли состояние)
task.Status = storage.StatusRunning
if err := w.Store.UpdateTask(ctx, task); 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/<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.notifyStatus(ctx, task, storage.StatusTimeout)
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.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 с объяснением.
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.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
}
task.Status = storage.StatusSuccess
if e := w.Store.UpdateTask(ctx, task); 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 с объяснением.
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.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) {
task.Status = storage.StatusFailed
if e := w.Store.UpdateTask(ctx, task); 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/<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
}