worker: add reviewer stage (branch→review→push) with R1-R6 errors
All checks were successful
CI / test (push) Successful in 47s
CI / build-and-package (amd64, linux) (push) Successful in 43s
CI / build-and-package (amd64, windows) (push) Successful in 40s

This commit is contained in:
Hermes
2026-08-16 21:21:21 +05:00
parent 130dc59840
commit 4a6ed57971
8 changed files with 610 additions and 65 deletions

View File

@@ -124,7 +124,9 @@ func (w *Worker) pollAndDispatch(ctx context.Context) error {
return nil
}
// runTask выполняет одну задачу: готовит репозитории, затем dev-агент через opencode.
// 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 {
@@ -155,68 +157,164 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
return fmt.Errorf("%w: %v", ErrClone, err)
}
// 3. рендерим промпт
prompt, err := RenderDevPrompt(DevPromptData{
Title: task.Title,
Goal: task.Goal,
Repos: repos,
Why: task.Why,
AC: task.AC,
})
if err != nil {
return fmt.Errorf("%w: render prompt: %v", ErrTrace, 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)
}
}
// 4. создаём трассу
trace := &storage.Trace{
TaskID: task.ID,
Agent: w.Agent,
Prompt: prompt,
}
traceID, err := w.Store.AppendTrace(ctx, trace)
if err != nil {
return fmt.Errorf("%w: create: %v", ErrTrace, err)
}
// 5. cwd — общий каталог (вариант A: один dev видит все репозитории).
// 3. cwd — общий каталог (вариант A: один dev видит все репозитории).
cwd := w.Worktree
// 6. запускаем dev-агент
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)
}
// Цикл dev → review, до maxReviewIterations.
var feedback []string
// 6b. сохраняем session_id из результата
if res.SessionID != "" {
_ = w.Store.UpdateTraceSessionID(ctx, traceID, res.SessionID)
}
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)
}
// 7. определяем результат по RC
output := res.Stdout
var traceStatus storage.TraceStatus
// 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)
}
switch {
case res.RC == 0:
task.Status = storage.StatusSuccess
traceStatus = storage.TraceSuccess
case res.RC == -1:
task.Status = storage.StatusTimeout
traceStatus = storage.TraceTimeout
default:
// 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
traceStatus = storage.TraceFailed
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
}
}
// 8. сохраняем результат
if e := w.Store.UpdateTask(ctx, task); e != nil {
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
// 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
}
w.finalizeTrace(ctx, traceID, traceStatus, output)
return nil
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.