feat: авто-уведомления о статусах и хендоффах задачи
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

This commit is contained in:
ki.sagidullin
2026-08-17 23:32:26 +05:00
parent 43266ea04f
commit 07bae203c3
4 changed files with 372 additions and 4 deletions

View File

@@ -22,6 +22,14 @@ type OpenCodeRunner interface {
// 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
@@ -39,6 +47,10 @@ type Worker struct {
// Через него Runner пишет live-шаги задачи; nil — наблюдение выключено.
Live *opencode.LiveRegistry
// Notify — нотифаер авто-уведомлений владельцу задачи (статусы + хендоффы
// dev↔reviewer). nil — уведомления выключены.
Notify Notifier
sem chan struct{} // семафор
cancel context.CancelFunc
@@ -55,6 +67,27 @@ func (w *Worker) runCtx(ctx context.Context, taskID int64) context.Context {
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 == "" {
@@ -163,6 +196,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
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 {
@@ -231,6 +265,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
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:
@@ -238,6 +273,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
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
}
@@ -245,6 +281,9 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
// 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 {
@@ -275,6 +314,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
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
}
@@ -289,11 +329,13 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
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
}
@@ -303,6 +345,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
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 {
@@ -330,12 +373,13 @@ func (w *Worker) reviewWithRetry(ctx context.Context, taskID int64, cwd, prompt
return v2, out2, tid2, nil
}
// failTask помечает задачу failed.
// 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 и статус трассы.

View File

@@ -8,6 +8,7 @@ import (
"os"
"os/exec"
"path/filepath"
"reflect"
"strconv"
"strings"
"testing"
@@ -17,6 +18,32 @@ import (
"github.com/kamelion/ratatoskr-go/internal/storage"
)
// fakeNotifier — фейковый нотифаер, собирающий все авто-уведомления воркера.
type fakeNotifier struct {
notifs []notifCall
}
// notifCall — одно перехваченное уведомление.
type notifCall struct {
taskID int64
chatID string
text string
}
func (f *fakeNotifier) Notify(_ context.Context, taskID int64, chatID, text string) error {
f.notifs = append(f.notifs, notifCall{taskID: taskID, chatID: chatID, text: text})
return nil
}
// notifTexts возвращает тексты уведомлений в порядке отправки.
func notifTexts(n *fakeNotifier) []string {
texts := make([]string, len(n.notifs))
for i, c := range n.notifs {
texts[i] = c.text
}
return texts
}
type mockRunnerWorker struct {
result *opencode.Result
err error
@@ -263,6 +290,186 @@ func TestWorkerReviewMaxIterations(t *testing.T) {
}
}
// TestWorkerStatusNotifications — happy path: владелец получает уведомления
// на каждый переход статуса со стороны воркера (running → success)
// и на хендофф dev→reviewer.
func TestWorkerStatusNotifications(t *testing.T) {
s := setupWorkerDB(t)
task := createReadyTask(t, s, "notif-ok")
n := &fakeNotifier{}
w := &Worker{
Store: s,
Runner: &mockRunnerWorker{result: &opencode.Result{RC: 0, Stdout: "done", SessionID: "sess-1"}},
Worktree: t.TempDir(),
Agent: "dev",
Notify: n,
}
seedFakeRepo(t, w.Worktree, "notif-ok")
ctx := context.Background()
if err := w.runTask(ctx, task); err != nil {
t.Fatalf("runTask: %v", err)
}
if len(n.notifs) != 3 {
t.Fatalf("уведомлений = %d, want 3 (running, dev→reviewer, success)", len(n.notifs))
}
prefix := "Задача #" + strconv.FormatInt(task.ID, 10)
want := []string{
prefix + ": running",
prefix + ": dev → reviewer (итерация 1)",
prefix + ": success",
}
if got := notifTexts(n); !reflect.DeepEqual(got, want) {
t.Errorf("уведомления = %#v, want %#v", got, want)
}
// все уведомления уходят владельцу задачи (task.ChatID)
for _, c := range n.notifs {
if c.chatID != task.ChatID {
t.Errorf("уведомление ушло в %q, want %q", c.chatID, task.ChatID)
}
if c.taskID != task.ID {
t.Errorf("уведомление для задачи %d, want %d", c.taskID, task.ID)
}
}
}
// TestWorkerHandoffNotifications — цикл dev↔review: уведомления на оба хендоффа
// (dev→reviewer и reviewer→dev на доработку) с номером задачи и итерации.
func TestWorkerHandoffNotifications(t *testing.T) {
s := setupWorkerDB(t)
task := createReadyTask(t, s, "notif-loop")
n := &fakeNotifier{}
w := &Worker{
Store: s,
Runner: &mockRunnerWorker{
result: &opencode.Result{RC: 0, Stdout: "done", SessionID: "sess-1"},
reviewSequence: []*opencode.Result{
reviewFailedRunner(),
{RC: 0, Stdout: `{"passed":true,"comments":[]}`},
},
},
Worktree: t.TempDir(),
Agent: "dev",
Notify: n,
}
seedFakeRepo(t, w.Worktree, "notif-loop")
ctx := context.Background()
if err := w.runTask(ctx, task); err != nil {
t.Fatalf("runTask: %v", err)
}
prefix := "Задача #" + strconv.FormatInt(task.ID, 10)
want := []string{
prefix + ": running",
prefix + ": dev → reviewer (итерация 1)",
prefix + ": reviewer → dev на доработку (итерация 1)",
prefix + ": dev → reviewer (итерация 2)",
prefix + ": success",
}
if got := notifTexts(n); !reflect.DeepEqual(got, want) {
t.Errorf("уведомления = %#v, want %#v", got, want)
}
}
// TestWorkerIterationsLimitNotification — при исчерпании лимита итераций
// владельцу уходит одно уведомление о failed (без дублей с running→failed).
func TestWorkerIterationsLimitNotification(t *testing.T) {
s := setupWorkerDB(t)
task := createReadyTask(t, s, "notif-lim")
n := &fakeNotifier{}
w := &Worker{
Store: s,
Runner: &mockRunnerWorker{
result: &opencode.Result{RC: 0, Stdout: "done", SessionID: "sess-1"},
reviewResult: reviewFailedRunner(),
},
Worktree: t.TempDir(),
Agent: "dev",
Notify: n,
}
seedFakeRepo(t, w.Worktree, "notif-lim")
ctx := context.Background()
if err := w.runTask(ctx, task); err != nil {
t.Fatalf("runTask: %v", err)
}
texts := notifTexts(n)
if len(texts) == 0 {
t.Fatal("нет уведомлений")
}
last := texts[len(texts)-1]
if !strings.Contains(last, "failed") {
t.Errorf("последнее уведомление = %q, want упоминание failed", last)
}
if !strings.Contains(last, "итераци") {
t.Errorf("последнее уведомление = %q, want упоминание лимита итераций", last)
}
// ровно одно уведомление о failed (running→failed не задваивается)
var failedCount int
for _, txt := range texts {
if strings.Contains(txt, ": failed") {
failedCount++
}
}
if failedCount != 1 {
t.Errorf("уведомлений о failed = %d, want ровно 1: %#v", failedCount, texts)
}
}
// TestWorkerTimeoutNotification — RC=-1 (таймаут dev) → уведомление о timeout.
func TestWorkerTimeoutNotification(t *testing.T) {
s := setupWorkerDB(t)
task := createReadyTask(t, s, "notif-timeout")
n := &fakeNotifier{}
w := &Worker{
Store: s,
Runner: &mockRunnerWorker{result: &opencode.Result{RC: -1, Stdout: ""}},
Worktree: t.TempDir(),
Notify: n,
}
seedFakeRepo(t, w.Worktree, "notif-timeout")
ctx := context.Background()
_ = w.runTask(ctx, task)
prefix := "Задача #" + strconv.FormatInt(task.ID, 10)
want := []string{prefix + ": running", prefix + ": timeout"}
if got := notifTexts(n); !reflect.DeepEqual(got, want) {
t.Errorf("уведомления = %#v, want %#v", got, want)
}
}
// TestWorkerSpawnErrorNotification — сбой запуска dev → уведомление о failed.
func TestWorkerSpawnErrorNotification(t *testing.T) {
s := setupWorkerDB(t)
task := createReadyTask(t, s, "notif-spawn")
n := &fakeNotifier{}
w := &Worker{
Store: s,
Runner: &mockRunnerWorker{err: errors.New("opencode not found")},
Worktree: t.TempDir(),
Notify: n,
}
seedFakeRepo(t, w.Worktree, "notif-spawn")
ctx := context.Background()
_ = w.runTask(ctx, task)
prefix := "Задача #" + strconv.FormatInt(task.ID, 10)
want := []string{prefix + ": running", prefix + ": failed"}
if got := notifTexts(n); !reflect.DeepEqual(got, want) {
t.Errorf("уведомления = %#v, want %#v", got, want)
}
}
// reviewNDJSONRunner возвращает вердикт ревьюера как реальный NDJSON-поток opencode,
// где JSON находится внутри последнего text-парта.
func reviewNDJSONRunner(v *reviewVerdict) *opencode.Result {