Files
ratatoskr-go/internal/worker/worker.go
Hermes bd3d825738
All checks were successful
CI / test (push) Successful in 46s
CI / build-and-package (amd64, darwin) (push) Successful in 36s
CI / build-and-package (amd64, linux) (push) Successful in 37s
CI / build-and-package (amd64, windows) (push) Successful in 44s
CI / build-and-package (arm64, darwin) (push) Successful in 35s
CI / build-and-package (arm64, linux) (push) Successful in 36s
feat: множественные репозитории (Repos) и клонирование в воркере
- Task.Repos []string (XML-колонка repos, обратная совместимость с repo)
- config: блок git {base_url, token}
- аналитик: ответ repos[], шаблон показывает список
- core: propose без repos → возврат в сбор (E1)
- worker вариант A: один dev из общего cwd, prepareRepos клонирует
  недостающие репо (git clone), validateRepoName (E3), ErrRepoNotGit (E4)
- ошибки E1-E4 в worker/errors.go
2026-08-16 09:18:15 +05:00

299 lines
9.0 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 выполняет одну задачу: готовит репозитории, затем dev-агент через opencode.
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)
}
// 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)
}
// 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 видит все репозитории).
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)
}
// 6b. сохраняем session_id из результата
if res.SessionID != "" {
_ = w.Store.UpdateTraceSessionID(ctx, traceID, res.SessionID)
}
// 7. определяем результат по RC
output := res.Stdout
var traceStatus storage.TraceStatus
switch {
case res.RC == 0:
task.Status = storage.StatusSuccess
traceStatus = storage.TraceSuccess
case res.RC == -1:
task.Status = storage.StatusTimeout
traceStatus = storage.TraceTimeout
default:
task.Status = storage.StatusFailed
traceStatus = storage.TraceFailed
}
// 8. сохраняем результат
if e := w.Store.UpdateTask(ctx, task); e != nil {
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
}
w.finalizeTrace(ctx, traceID, traceStatus, output)
return 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 {
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
}