From 46410903985acc182fc291fe939fd3308e604228 Mon Sep 17 00:00:00 2001 From: Hermes Date: Sat, 15 Aug 2026 09:14:34 +0500 Subject: [PATCH] =?UTF-8?q?worker:=20=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2=D0=BB?= =?UTF-8?q?=D0=B5=D0=BD=20=D0=BF=D0=B0=D0=BA=D0=B5=D1=82=20worker=20(W1-W5?= =?UTF-8?q?,=20poll-=D1=86=D0=B8=D0=BA=D0=BB,=20runTask,=20=D0=BF=D1=80?= =?UTF-8?q?=D0=BE=D0=BC=D0=BF=D1=82,=20=D1=82=D0=B5=D1=81=D1=82=D1=8B)=20+?= =?UTF-8?q?=20SetMaxOpenConns(1)=20=D0=B4=D0=BB=D1=8F=20SQLITE=5FBUSY?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/storage/storage.go | 3 + internal/storage/traces.go | 9 + internal/worker/errors.go | 21 +++ internal/worker/prompt.go | 42 +++++ internal/worker/worker.go | 218 +++++++++++++++++++++ internal/worker/worker_test.go | 334 +++++++++++++++++++++++++++++++++ 6 files changed, 627 insertions(+) create mode 100644 internal/worker/errors.go create mode 100644 internal/worker/prompt.go create mode 100644 internal/worker/worker.go create mode 100644 internal/worker/worker_test.go diff --git a/internal/storage/storage.go b/internal/storage/storage.go index 3196944..93c9304 100644 --- a/internal/storage/storage.go +++ b/internal/storage/storage.go @@ -81,6 +81,9 @@ func Open(ctx context.Context, path string) (*Storage, error) { if err != nil { return nil, fmt.Errorf("%w: open: %w", ErrDB, err) } + // Один connection — сериализует доступ, исключает SQLITE_BUSY + // при конкурентной записи из нескольких горутин (worker). + db.SetMaxOpenConns(1) // Прагмы: WAL + синхронность pragmas := []string{ "PRAGMA journal_mode=WAL", diff --git a/internal/storage/traces.go b/internal/storage/traces.go index 7dd2abd..549fe91 100644 --- a/internal/storage/traces.go +++ b/internal/storage/traces.go @@ -53,6 +53,15 @@ func (s *Storage) UpdateTraceOutput(ctx context.Context, traceID int64, output s return nil } +// UpdateTraceSessionID обновляет session_id трассы (например, после запуска opencode). +func (s *Storage) UpdateTraceSessionID(ctx context.Context, traceID int64, sessionID string) error { + _, err := s.db.ExecContext(ctx, `UPDATE traces SET session_id=? WHERE id=?`, sessionID, traceID) + if err != nil { + return fmt.Errorf("%w: update trace session_id %d: %w", ErrDB, traceID, err) + } + return nil +} + // GetTraces возвращает трассы для задачи, отсортированные по started_at. func (s *Storage) GetTraces(ctx context.Context, taskID int64) ([]*Trace, error) { rows, err := s.db.QueryContext(ctx, ` diff --git a/internal/worker/errors.go b/internal/worker/errors.go new file mode 100644 index 0000000..ef76f27 --- /dev/null +++ b/internal/worker/errors.go @@ -0,0 +1,21 @@ +package worker + +import "errors" + +// Классы ошибок W1–W5. +var ( + // W1 — ошибка опроса БД. + ErrPoll = errors.New("W1: poll error") + + // W2 — задача в невалидном статусе. + ErrLaunch = errors.New("W2: launch error") + + // W3 — ошибка создания/обновления трассы. + ErrTrace = errors.New("W3: trace error") + + // W4 — ошибка обновления статуса задачи. + ErrUpdate = errors.New("W4: update error") + + // W5 — превышена параллельность, задача пропущена. + ErrConcurrencyLimit = errors.New("W5: concurrency limit") +) \ No newline at end of file diff --git a/internal/worker/prompt.go b/internal/worker/prompt.go new file mode 100644 index 0000000..dafe425 --- /dev/null +++ b/internal/worker/prompt.go @@ -0,0 +1,42 @@ +package worker + +import ( + "strings" + "text/template" +) + +// devPromptTemplate — промпт для dev-агента при запуске задачи. +var devPromptTemplate = template.Must(template.New("dev").Parse(`Ты — dev-агент, реализуешь задачу в репозитории. + +**Задача:** +{{if .Title}}Название: {{.Title}}{{end}} +{{if .Goal}}Цель: {{.Goal}}{{end}} +{{if .Why}}Зачем: {{.Why}}{{end}} +{{if .Repo}}Репозиторий: {{.Repo}}{{end}} +{{if .AC}}Критерии готовности: +{{.AC}}{{end}} + +**Инструкции:** +1. Напиши код, реализующий задачу. +2. Убедись, что все acceptance criteria выполнены. +3. В процессе работы пользуйся встроенными инструментами opencode (чтение файлов, поиск, редактирование). +4. По окончании работы верни краткий отчёт о том, что сделано. +`)) + +// DevPromptData — данные для рендера dev-промпта. +type DevPromptData struct { + Title string + Goal string + Repo string + Why string + AC string +} + +// RenderDevPrompt собирает промпт для dev-агента. +func RenderDevPrompt(data DevPromptData) (string, error) { + var buf strings.Builder + if err := devPromptTemplate.Execute(&buf, data); err != nil { + return "", err + } + return buf.String(), nil +} \ No newline at end of file diff --git a/internal/worker/worker.go b/internal/worker/worker.go new file mode 100644 index 0000000..5c63140 --- /dev/null +++ b/internal/worker/worker.go @@ -0,0 +1,218 @@ +package worker + +import ( + "context" + "fmt" + "log" + "path/filepath" + "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 // базовый путь, task.Repo — относительно него + Agent string // default "dev" + Interval time.Duration // интервал опроса БД + MaxJobs int // макс. параллельных задач + + 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) + } + + // 2. ставим running + task.Status = storage.StatusRunning + if err := w.Store.UpdateTask(ctx, task); err != nil { + return fmt.Errorf("%w: set running: %v", ErrUpdate, err) + } + + // 3. рендерим промпт + prompt, err := RenderDevPrompt(DevPromptData{ + Title: task.Title, + Goal: task.Goal, + Repo: task.Repo, + 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 + cwd := w.resolveCwd(task.Repo) + + // 6. запускаем dev-агент + res, resErr := w.Runner.Run(ctx, prompt, cwd, w.Agent, "") + if resErr != nil { + // O1 ErrSpawn — не смог запустить бинарь + task.Status = storage.StatusFailed + if e := w.Store.UpdateTask(ctx, task); e != nil { + err = fmt.Errorf("%w: set failed: %v", ErrUpdate, e) + return + } + w.finalizeTrace(ctx, traceID, storage.TraceFailed, resErr.Error()) + err = fmt.Errorf("%w: spawn: %v", ErrLaunch, resErr) + return + } + + // 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 { + err = fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e) + return + } + w.finalizeTrace(ctx, traceID, traceStatus, output) + return nil +} + +// 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) + } +} + +func (w *Worker) resolveCwd(repo string) string { + if repo == "" { + return w.Worktree + } + return filepath.Join(w.Worktree, repo) +} \ No newline at end of file diff --git a/internal/worker/worker_test.go b/internal/worker/worker_test.go new file mode 100644 index 0000000..732e851 --- /dev/null +++ b/internal/worker/worker_test.go @@ -0,0 +1,334 @@ +package worker + +import ( + "context" + "errors" + "fmt" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/kamelion/ratatoskr-go/internal/opencode" + "github.com/kamelion/ratatoskr-go/internal/storage" +) + +type mockRunnerWorker struct { + result *opencode.Result + err error +} + +func (m *mockRunnerWorker) Run(_ context.Context, _, _, _, _ string) (*opencode.Result, error) { + return m.result, m.err +} + +func setupWorkerDB(t *testing.T) *storage.Storage { + t.Helper() + ctx := context.Background() + f := filepath.Join(t.TempDir(), "test.db") + s, err := storage.Open(ctx, f) + if err != nil { + t.Fatalf("open storage: %v", err) + } + t.Cleanup(func() { s.Close() }) + return s +} + +func createReadyTask(t *testing.T, s *storage.Storage, title string) *storage.Task { + t.Helper() + ctx := context.Background() + task := &storage.Task{ + ChatID: "tg://worker-test", + Title: title, + Goal: "сделать " + title, + Repo: "test/" + title, + Why: "для теста", + AC: "работает", + TaskTag: "test-" + title, + } + id, err := s.CreateTask(ctx, task) + if err != nil { + t.Fatalf("create task: %v", err) + } + task.ID = id + task.Status = storage.StatusCollecting + if err := s.UpdateTask(ctx, task); err != nil { + t.Fatalf("set collecting: %v", err) + } + task.Status = storage.StatusReady + if err := s.UpdateTask(ctx, task); err != nil { + t.Fatalf("set ready: %v", err) + } + task, _ = s.GetTask(ctx, id) + return task +} + +func TestWorkerHappyPath(t *testing.T) { + s := setupWorkerDB(t) + task := createReadyTask(t, s, "calc") + + w := &Worker{ + Store: s, + Runner: &mockRunnerWorker{result: &opencode.Result{RC: 0, Stdout: "done", SessionID: "sess-1"}}, + Worktree: t.TempDir(), + Agent: "dev", + } + + ctx := context.Background() + if err := w.runTask(ctx, task); err != nil { + t.Fatalf("runTask: %v", err) + } + + task, err := s.GetTask(ctx, task.ID) + if err != nil { + t.Fatalf("get task: %v", err) + } + if task.Status != storage.StatusSuccess { + t.Errorf("status = %q, want success", task.Status) + } + + traces, err := s.GetTraces(ctx, task.ID) + if err != nil { + t.Fatalf("get traces: %v", err) + } + if len(traces) != 1 { + t.Fatalf("got %d traces, want 1", len(traces)) + } + if traces[0].Status != storage.TraceSuccess { + t.Errorf("trace status = %q, want success", traces[0].Status) + } + if traces[0].Agent != "dev" { + t.Errorf("agent = %q, want dev", traces[0].Agent) + } + if traces[0].SessionID != "sess-1" { + t.Errorf("session = %q, want sess-1", traces[0].SessionID) + } +} + +func TestWorkerTimeout(t *testing.T) { + s := setupWorkerDB(t) + task := createReadyTask(t, s, "slow") + + w := &Worker{ + Store: s, + Runner: &mockRunnerWorker{result: &opencode.Result{RC: -1, Stdout: ""}}, + Worktree: t.TempDir(), + } + + ctx := context.Background() + _ = w.runTask(ctx, task) + + task, err := s.GetTask(ctx, task.ID) + if err != nil { + t.Fatalf("get task: %v", err) + } + if task.Status != storage.StatusTimeout { + t.Errorf("status = %q, want timeout", task.Status) + } + + traces, err := s.GetTraces(ctx, task.ID) + if err != nil { + t.Fatalf("get traces: %v", err) + } + if len(traces) != 1 { + t.Fatalf("got %d traces, want 1", len(traces)) + } + if traces[0].Status != storage.TraceTimeout { + t.Errorf("trace status = %q, want timeout", traces[0].Status) + } +} + +func TestWorkerSpawnError(t *testing.T) { + s := setupWorkerDB(t) + task := createReadyTask(t, s, "spawn-fail") + + w := &Worker{ + Store: s, + Runner: &mockRunnerWorker{err: errors.New("opencode not found")}, + Worktree: t.TempDir(), + } + + ctx := context.Background() + _ = w.runTask(ctx, task) + + task, err := s.GetTask(ctx, task.ID) + if err != nil { + t.Fatalf("get task: %v", err) + } + if task.Status != storage.StatusFailed { + t.Errorf("status = %q, want failed", task.Status) + } + + traces, err := s.GetTraces(ctx, task.ID) + if err != nil { + t.Fatalf("get traces: %v", err) + } + if len(traces) != 1 { + t.Fatalf("got %d traces, want 1", len(traces)) + } + if traces[0].Status != storage.TraceFailed { + t.Errorf("trace status = %q, want failed", traces[0].Status) + } +} + +func TestWorkerNonZeroExit(t *testing.T) { + s := setupWorkerDB(t) + task := createReadyTask(t, s, "fail") + + w := &Worker{ + Store: s, + Runner: &mockRunnerWorker{result: &opencode.Result{RC: 7, Stdout: "error"}}, + Worktree: t.TempDir(), + } + + ctx := context.Background() + _ = w.runTask(ctx, task) + + task, err := s.GetTask(ctx, task.ID) + if err != nil { + t.Fatalf("get task: %v", err) + } + if task.Status != storage.StatusFailed { + t.Errorf("status = %q, want failed", task.Status) + } + + traces, err := s.GetTraces(ctx, task.ID) + if err != nil { + t.Fatalf("get traces: %v", err) + } + if len(traces) != 1 { + t.Fatalf("got %d traces, want 1", len(traces)) + } + if traces[0].Status != storage.TraceFailed { + t.Errorf("trace status = %q, want failed", traces[0].Status) + } + if traces[0].Output != "error" { + t.Errorf("output = %q, want error", traces[0].Output) + } +} + +func TestWorkerBadStatus(t *testing.T) { + s := setupWorkerDB(t) + ctx := context.Background() + + task := &storage.Task{ChatID: "tg://bad", Title: "bad-status", TaskTag: "bad"} + id, err := s.CreateTask(ctx, task) + if err != nil { + t.Fatalf("create task: %v", err) + } + task, _ = s.GetTask(ctx, id) + + err = (&Worker{Store: s}).runTask(ctx, task) + if err == nil { + t.Fatal("expected error for non-ready task") + } + if !errors.Is(err, ErrLaunch) { + t.Errorf("err = %v, want W2", err) + } +} + +func TestWorkerPromptRendered(t *testing.T) { + s := setupWorkerDB(t) + task := createReadyTask(t, s, "prompt-test") + + w := &Worker{ + Store: s, + Runner: &mockRunnerWorker{result: &opencode.Result{RC: 0, Stdout: "ok", SessionID: "s"}}, + Worktree: t.TempDir(), + } + + ctx := context.Background() + _ = w.runTask(ctx, task) + + traces, err := s.GetTraces(ctx, task.ID) + if err != nil { + t.Fatalf("get traces: %v", err) + } + if len(traces) == 0 { + t.Fatal("no traces") + } + tr := traces[0] + if tr.Prompt == "" { + t.Error("prompt empty — шаблон не срендерился") + } + if !strings.Contains(tr.Prompt, "prompt-test") { + t.Error("prompt не содержит название задачи") + } + if !strings.Contains(tr.Prompt, "test/prompt-test") { + t.Error("prompt не содержит repo") + } +} + +func TestWorkerResolveCwd(t *testing.T) { + base := "/opt/data/src" + w := &Worker{Worktree: base} + + if got := w.resolveCwd(""); got != base { + t.Errorf("empty repo: got %q, want %q", got, base) + } + if got := w.resolveCwd("tools/calc"); got != filepath.Join(base, "tools/calc") { + t.Errorf("repo: got %q, want %q", got, filepath.Join(base, "tools/calc")) + } +} + +func TestWorkerStartStop(t *testing.T) { + s := setupWorkerDB(t) + w := &Worker{ + Store: s, + Runner: &mockRunnerWorker{result: &opencode.Result{RC: 0}}, + Worktree: t.TempDir(), + MaxJobs: 1, + Interval: 50 * time.Millisecond, + } + + ctx := context.Background() + w.Start(ctx) + + time.Sleep(150 * time.Millisecond) + w.Stop() +} + +func TestWorkerSemaphore(t *testing.T) { + s := setupWorkerDB(t) + + // создаём 2 ready-задачи + for i := 0; i < 2; i++ { + createReadyTask(t, s, fmt.Sprintf("task-%d", i)) + } + + w := &Worker{ + Store: s, + Runner: &mockRunnerWorker{result: &opencode.Result{RC: 0, Stdout: "ok"}}, + Worktree: t.TempDir(), + MaxJobs: 1, + Interval: 50 * time.Millisecond, + } + w.sem = make(chan struct{}, 1) + w.sem <- struct{}{} + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + + // первый poll — запустит 1 задачу (макс. 1) + w.pollAndDispatch(ctx) + time.Sleep(200 * time.Millisecond) + + // 1 должна быть success, 1 — всё ещё ready + success, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess}) + ready, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusReady}) + if len(success) != 1 { + t.Errorf("success = %d, want 1 (ready=%d)", len(success), len(ready)) + } + if len(ready) != 1 { + t.Errorf("ready = %d, want 1", len(ready)) + } + + // первая завершилась и вернула токен в сем — можем диспатчить вторую + w.pollAndDispatch(ctx) + time.Sleep(200 * time.Millisecond) + + success, _ = s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess}) + if len(success) != 2 { + t.Errorf("после освобождения слота success = %d, want 2", len(success)) + } +} \ No newline at end of file