diff --git a/internal/app/app.go b/internal/app/app.go index 537db76..5b330b4 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -154,6 +154,7 @@ func New(configPath, version, updateToken string) (*App, error) { GitBaseURL: cfg.Git.BaseURL, GitToken: cfg.Git.Token, Live: live, + Notify: a, // авто-уведомления владельцу задачи через Router } a.Router = router a.Worker = w @@ -326,6 +327,16 @@ func (a *App) send(ctx context.Context, uid chat.UserID, text string) { } } +// Notify реализует worker.Notifier: авто-уведомление владельцу задачи через +// chat.Router.Send (переходы статусов и хендоффы dev↔reviewer со стороны воркера). +// Router nil (тесты без Router / ранняя инициализация) — тихо пропускаем. +func (a *App) Notify(ctx context.Context, taskID int64, chatID, text string) error { + if a.Router == nil { + return nil + } + return a.Router.Send(ctx, chat.UserID(chatID), chat.Message{Text: text}) +} + // cmdName извлекает команду (первое слово до пробела, нижний регистр). func cmdName(text string) string { s := strings.TrimSpace(text) diff --git a/internal/app/e2e_test.go b/internal/app/e2e_test.go index 071b3bb..bc9548e 100644 --- a/internal/app/e2e_test.go +++ b/internal/app/e2e_test.go @@ -18,6 +18,7 @@ import ( "os" "os/exec" "path/filepath" + "reflect" "strings" "testing" "time" @@ -129,6 +130,7 @@ func e2eAssemble(t *testing.T) (*App, string, *e2eChannel) { Interval: 30 * time.Millisecond, MaxJobs: 1, Live: opencode.NewLiveRegistry(), + Notify: a, // авто-уведомления владельцу через Router (как в app.New) // GitToken зададим пустым: origin в seed-репо локальный (file path), // http.extraHeader не нужен для локального пуша. } @@ -139,13 +141,13 @@ func e2eAssemble(t *testing.T) (*App, string, *e2eChannel) { // e2eChannel — минимальный fake-канал для перехвата исходящих // и доставки входящих через роутер (как реальный канал). type e2eChannel struct { - onMsg func(chat.Incoming) + onMsg chat.Handler sent []chat.Message } func (c *e2eChannel) Run(_ context.Context) error { return nil } -func (c *e2eChannel) OnMessage(f func(chat.Incoming)) { - c.onMsg = f +func (c *e2eChannel) OnMessage(h chat.Handler) { + c.onMsg = h } func (c *e2eChannel) Send(_ context.Context, _ chat.Address, m chat.Message) error { c.sent = append(c.sent, m) @@ -719,4 +721,108 @@ func TestE2EWorkerDoesNotTakeUnconfirmed(t *testing.T) { } t.Logf("WORKER-CONSENT OK: в ready воркер не трогал, после «создавай» → success, ветка %s", branch) +} + +// TestE2ENotificationsOnTransitions — авто-уведомления владельцу на каждый +// переход статуса задачи и отсутствие задвоения на approved. +// +// Проверяет сквозной поток сообщений через chat.Router: +// - ready → approved («создавай»): ровно ОДНО сообщение (ответ «одобрена»), +// отдельное уведомление "Задача #N: approved" НЕ дублируется; +// - со стороны воркера: уведомления running → dev→reviewer → success +// доходят владельцу через Router.Send. +func TestE2ENotificationsOnTransitions(t *testing.T) { + a, worktree, fake := e2eAssemble(t) + a.seedFakeRepo(t, worktree, "calc") + + ctx := context.Background() + uid := chat.UserID("u-notif") + + // --- 1. постановка: /start → draft→collecting (приветствие) --- + fake.deliver(uid, "/start") + task, err := a.Store.GetActiveTaskByChatID(ctx, string(uid)) + if err != nil { + t.Fatalf("get task after /start: %v", err) + } + taskID := task.ID + + // --- 2. текст → collecting→ready (сводка черновика) --- + fake.deliver(uid, "Сделай калькулятор в calc") + task, err = a.Store.GetTask(ctx, taskID) + if err != nil { + t.Fatalf("get task: %v", err) + } + if task.Status != storage.StatusReady { + t.Fatalf("status = %q, want ready", task.Status) + } + + // --- 3. «создавай» → ready→approved: один ответ, без дубля-уведомления --- + fake.deliver(uid, "создавай") + task, err = a.Store.GetTask(ctx, taskID) + if err != nil { + t.Fatalf("get task after создавай: %v", err) + } + if task.Status != storage.StatusApproved { + t.Fatalf("status = %q, want approved", task.Status) + } + + approvedReplies := 0 + var approvedText string + approvedNotifs := 0 + for _, m := range fake.sent { + if strings.Contains(m.Text, "одобрена") { + approvedReplies++ + approvedText = m.Text + } + if strings.Contains(m.Text, fmt.Sprintf("Задача #%d: approved", taskID)) { + approvedNotifs++ + } + } + if approvedReplies != 1 { + t.Errorf("сообщений об одобрении = %d, want ровно 1 (нет задвоения)", approvedReplies) + } + if approvedNotifs != 0 { + t.Errorf("отдельное уведомление 'Задача #%d: approved' продублировано (%d раз)", taskID, approvedNotifs) + } + if !strings.Contains(approvedText, "одобрена") { + t.Errorf("ответ об одобрении = %q", approvedText) + } + + // --- 4. воркер: running → dev→reviewer → success через Router --- + wkCtx, wkCancel := context.WithCancel(ctx) + defer wkCancel() + a.Worker.Start(wkCtx) + + deadline := time.Now().Add(30 * time.Second) + for { + tk, err := a.Store.GetTask(ctx, taskID) + if err != nil { + t.Fatalf("get task: %v", err) + } + if tk.Status == storage.StatusSuccess || tk.Status == storage.StatusFailed { + break + } + if time.Now().After(deadline) { + t.Fatalf("таймаут ожидания success, последний статус %q", tk.Status) + } + time.Sleep(50 * time.Millisecond) + } + + // финальные статус-уведомления воркера в правильном порядке + wantOrder := []string{ + fmt.Sprintf("Задача #%d: running", taskID), + fmt.Sprintf("Задача #%d: dev → reviewer (итерация 1)", taskID), + fmt.Sprintf("Задача #%d: success", taskID), + } + var gotOrder []string + for _, m := range fake.sent { + if strings.HasPrefix(m.Text, fmt.Sprintf("Задача #%d:", taskID)) { + gotOrder = append(gotOrder, m.Text) + } + } + if !reflect.DeepEqual(gotOrder, wantOrder) { + t.Errorf("уведомления воркера = %#v, want %#v", gotOrder, wantOrder) + } + + t.Logf("NOTIFICATIONS OK: задача #%d, сообщений одобрения=%d, уведомления=%v", taskID, approvedReplies, gotOrder) } \ No newline at end of file diff --git a/internal/worker/worker.go b/internal/worker/worker.go index 4779b64..8225337 100644 --- a/internal/worker/worker.go +++ b/internal/worker/worker.go @@ -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 и статус трассы. diff --git a/internal/worker/worker_test.go b/internal/worker/worker_test.go index db23919..43db6e9 100644 --- a/internal/worker/worker_test.go +++ b/internal/worker/worker_test.go @@ -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 {