From 4957a5552eb3524f59e1773441910d931aacac2c Mon Sep 17 00:00:00 2001 From: "ki.sagidullin" Date: Thu, 20 Aug 2026 04:30:29 +0500 Subject: [PATCH] =?UTF-8?q?feat(events):=20=D0=B2=D0=BD=D0=B5=D0=B4=D1=80?= =?UTF-8?q?=D0=B8=D1=82=D1=8C=20Publisher=20=D0=B2=20Core/Worker/Analyst?= =?UTF-8?q?=20+=20=D1=88=D0=B8=D0=BD=D1=8B=20=D0=B2=20app?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Core: setStatus публикует TaskStatusChanged на каждом переходе (handleStart/Cancel/Retry/Consent/Edit/Turn, runDecide). - Worker: setStatus публикует TaskStatusChanged (running/timeout/failed/ success, failTask), notify остаётся авто-уведомлением. - Analyst: Decide публикует AgentActivity (stage=decide). - app: создаёт доменную (256) и логовую (1024) шины, подключает их к Core/Worker/Analyst, редиректит log в обе (stderr + логовая шина). --- internal/analyst/analyst.go | 13 +++++++++ internal/app/app.go | 20 ++++++++++++++ internal/core/core.go | 54 ++++++++++++++++++++++++------------- internal/worker/worker.go | 51 ++++++++++++++++++++++------------- 4 files changed, 102 insertions(+), 36 deletions(-) diff --git a/internal/analyst/analyst.go b/internal/analyst/analyst.go index befe9ae..d7e61e4 100644 --- a/internal/analyst/analyst.go +++ b/internal/analyst/analyst.go @@ -8,6 +8,7 @@ import ( "strings" "github.com/kamelion/ratatoskr-go/internal/core" + "github.com/kamelion/ratatoskr-go/internal/events" "github.com/kamelion/ratatoskr-go/internal/opencode" "github.com/kamelion/ratatoskr-go/internal/storage" ) @@ -35,6 +36,16 @@ type Analyst struct { Runner OpenCodeRunner Worktree string // каталог, откуда запускать opencode run Agent string // имя агента (default "analyst") + + // Events — издатель доменных событий для UI. nil — события выключены. + Events events.Publisher +} + +// publish отправляет доменное событие, если задан издатель. +func (a *Analyst) publish(e events.Event) { + if a.Events != nil { + a.Events.Publish(e) + } } // AnalystResponse — структура JSON-ответа аналитика. @@ -58,6 +69,8 @@ func (a *Analyst) Decide(ctx context.Context, history []core.Message, draft stor agent = "analyst" } + a.publish(events.AgentActivity{TaskID: draft.ID, Agent: agent, Stage: "decide"}) + if len(history) == 0 && !force { return core.Decision{}, fmt.Errorf("%w: пустая история диалога", ErrNotReady) } diff --git a/internal/app/app.go b/internal/app/app.go index 93e77d1..de7ffdb 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -5,6 +5,7 @@ import ( "context" "errors" "fmt" + "io" "log" "os" "os/signal" @@ -19,6 +20,7 @@ import ( "github.com/kamelion/ratatoskr-go/internal/chat/telegram" "github.com/kamelion/ratatoskr-go/internal/config" "github.com/kamelion/ratatoskr-go/internal/core" + "github.com/kamelion/ratatoskr-go/internal/events" "github.com/kamelion/ratatoskr-go/internal/opencode" "github.com/kamelion/ratatoskr-go/internal/storage" "github.com/kamelion/ratatoskr-go/internal/update" @@ -61,6 +63,10 @@ type App struct { Updater *update.Updater tg *telegram.Channel // сохранена для Run pool *opencode.Pool // пул opencode serve-серверов (API-режим) + + // Events — доменная шина UI; LogEvents — шина логов (панель «Логи»). + Events *events.Bus + LogEvents *events.LogBus } // New читает конфиг и собирает все зависимости. @@ -144,6 +150,19 @@ func New(configPath, version, updateToken string) (*App, error) { CoreCtx: coreCtx, pool: ocPool, } + + // Шины событий для UI: доменная (статусы/история/трейсы) и логовая + // (сырые строки в панель «Логи»). Издатели подключаются к Core/Worker/ + // Analyst; на время headless (--noui) лог дублируется в обе шины, а UI + // просто подписывается на уже существующие. + a.Events = events.New(256) + a.LogEvents = events.NewLogBus(1024) + log.SetOutput(io.MultiWriter(os.Stderr, events.NewLogWriter(a.LogEvents))) + + // Подключаем издателей событий к Core/Worker/Analyst. + coreCtx.Events = a.Events + analystCtx.Events = a.Events + router := chat.NewRouter(a.handleIncoming) // Telegram-канал @@ -166,6 +185,7 @@ func New(configPath, version, updateToken string) (*App, error) { GitToken: cfg.Git.Token, Live: live, Notify: a, // авто-уведомления владельцу задачи через Router + Events: a.Events, } a.Router = router a.Worker = w diff --git a/internal/core/core.go b/internal/core/core.go index 3f53ad6..4da4306 100644 --- a/internal/core/core.go +++ b/internal/core/core.go @@ -5,6 +5,7 @@ import ( "fmt" "strings" + "github.com/kamelion/ratatoskr-go/internal/events" "github.com/kamelion/ratatoskr-go/internal/storage" ) @@ -18,6 +19,27 @@ type Core struct { MaxQuestionsPerTurn int // Live — опциональный просмотр живой сессии задачи (для /status N). Live LiveProber + // Events — издатель доменных событий для UI. nil — события выключены. + Events events.Publisher +} + +// publish отправляет доменное событие, если задан издатель. +func (c *Core) publish(e events.Event) { + if c.Events != nil { + c.Events.Publish(e) + } +} + +// setStatus переводит задачу в новый статус: сохраняет в БД и публикует событие. +// from фиксируется до перехода (для TaskStatusChanged). +func (c *Core) setStatus(ctx context.Context, task *storage.Task, to storage.Status) error { + from := task.Status + task.Status = to + if err := c.Store.UpdateTask(ctx, task); err != nil { + return err + } + c.publish(events.TaskStatusChanged{ID: task.ID, From: from, To: to}) + return nil } // LiveProber — абстракция за журналом живых сессий (реализация — *opencode.LiveRegistry). @@ -117,11 +139,10 @@ func (c *Core) handleStart(ctx context.Context, taskID int64) (Result, error) { if err != nil { return Result{}, err } - task.Status = storage.StatusCollecting - if err := c.Store.ClearHistory(ctx, task.ID); err != nil { + if err := c.setStatus(ctx, task, storage.StatusCollecting); err != nil { return Result{}, err } - if err := c.Store.UpdateTask(ctx, task); err != nil { + if err := c.Store.ClearHistory(ctx, task.ID); err != nil { return Result{}, err } return Result{ @@ -137,8 +158,7 @@ func (c *Core) handleCancel(ctx context.Context, taskID int64) (Result, error) { if err != nil { return Result{}, err } - task.Status = storage.StatusCancelled - if err := c.Store.UpdateTask(ctx, task); err != nil { + if err := c.setStatus(ctx, task, storage.StatusCancelled); err != nil { return Result{}, err } return Result{ @@ -202,11 +222,10 @@ func (c *Core) handleRetry(ctx context.Context, rest string) (Result, error) { Status: task.Status, }, nil } - task.Status = storage.StatusCollecting - if err := c.Store.ClearHistory(ctx, id); err != nil { + if err := c.setStatus(ctx, task, storage.StatusCollecting); err != nil { return Result{}, err } - if err := c.Store.UpdateTask(ctx, task); err != nil { + if err := c.Store.ClearHistory(ctx, id); err != nil { return Result{}, err } return Result{ @@ -272,8 +291,7 @@ func (c *Core) notFoundReply(ctx context.Context, id int64, err error) (Result, // handleConsent одобряет задачу: ready → approved (финальное одобрение, // после которого воркер забирает задачу на выполнение). func (c *Core) handleConsent(ctx context.Context, task *storage.Task) (Result, error) { - task.Status = storage.StatusApproved - if err := c.Store.UpdateTask(ctx, task); err != nil { + if err := c.setStatus(ctx, task, storage.StatusApproved); err != nil { return Result{}, err } return Result{ @@ -285,6 +303,7 @@ func (c *Core) handleConsent(ctx context.Context, task *storage.Task) (Result, e // handleEdit — правка черновика в фазе ready → снова сбор + аналитик. func (c *Core) handleEdit(ctx context.Context, task *storage.Task, text string) (Result, error) { + from := task.Status task.Status = storage.StatusCollecting if err := c.Store.AppendHistory(ctx, task.ID, "user", text); err != nil { return Result{}, err @@ -292,6 +311,7 @@ func (c *Core) handleEdit(ctx context.Context, task *storage.Task, text string) if err := c.Store.UpdateTask(ctx, task); err != nil { return Result{}, err } + c.publish(events.TaskStatusChanged{ID: task.ID, From: from, To: task.Status}) return c.runDecide(ctx, task, false) } @@ -300,10 +320,12 @@ func (c *Core) handleTurn(ctx context.Context, task *storage.Task, text string) // переход draft→collecting нужно персистить до вызова аналитика, // иначе propose сделает draft→ready (невалидно) if task.Status != storage.StatusCollecting { + from := task.Status task.Status = storage.StatusCollecting if err := c.Store.UpdateTask(ctx, task); err != nil { return Result{}, err } + c.publish(events.TaskStatusChanged{ID: task.ID, From: from, To: task.Status}) } if text != "" { if err := c.Store.AppendHistory(ctx, task.ID, "user", text); err != nil { @@ -348,8 +370,7 @@ func (c *Core) runDecide(ctx context.Context, task *storage.Task, force bool) (R switch decision.Phase { case "abort": - task.Status = storage.StatusAborted - if err := c.Store.UpdateTask(ctx, task); err != nil { + if err := c.setStatus(ctx, task, storage.StatusAborted); err != nil { return Result{}, err } reply := decision.ChatReply @@ -363,15 +384,13 @@ func (c *Core) runDecide(ctx context.Context, task *storage.Task, force bool) (R applyDraft(task, decision.Draft) // E1: propose/ready без репозиториев → остаёмся в сборе, просим уточнить. if len(task.EffectiveRepos()) == 0 { - task.Status = storage.StatusCollecting - if err := c.Store.UpdateTask(ctx, task); err != nil { + if err := c.setStatus(ctx, task, storage.StatusCollecting); err != nil { return Result{}, err } reply := decChatReply(decision, "Укажи, в каком репозитории(ях) вести работу.") return Result{Reply: reply, TaskID: task.ID, Status: task.Status}, nil } - task.Status = storage.StatusReady - if err := c.Store.UpdateTask(ctx, task); err != nil { + if err := c.setStatus(ctx, task, storage.StatusReady); err != nil { return Result{}, err } return Result{ @@ -382,8 +401,7 @@ func (c *Core) runDecide(ctx context.Context, task *storage.Task, force bool) (R default: // "ask" applyDraft(task, decision.Draft) - task.Status = storage.StatusCollecting - if err := c.Store.UpdateTask(ctx, task); err != nil { + if err := c.setStatus(ctx, task, storage.StatusCollecting); err != nil { return Result{}, err } reply := buildAskReply(decision, c.MaxQuestionsPerTurn) diff --git a/internal/worker/worker.go b/internal/worker/worker.go index 8225337..55b137d 100644 --- a/internal/worker/worker.go +++ b/internal/worker/worker.go @@ -10,6 +10,7 @@ import ( "strings" "time" + "github.com/kamelion/ratatoskr-go/internal/events" "github.com/kamelion/ratatoskr-go/internal/opencode" "github.com/kamelion/ratatoskr-go/internal/storage" ) @@ -56,6 +57,27 @@ type Worker struct { // подменяемый poll для тестов pollFn PollTaskFunc + + // Events — издатель доменных событий для UI. nil — события выключены. + Events events.Publisher +} + +// publish отправляет доменное событие, если задан издатель. +func (w *Worker) publish(e events.Event) { + if w.Events != nil { + w.Events.Publish(e) + } +} + +// setStatus переводит задачу в новый статус: сохраняет в БД и публикует событие. +func (w *Worker) setStatus(ctx context.Context, task *storage.Task, to storage.Status) error { + from := task.Status + task.Status = to + if err := w.Store.UpdateTask(ctx, task); err != nil { + return err + } + w.publish(events.TaskStatusChanged{ID: task.ID, From: from, To: to}) + return nil } // runCtx оборачивает контекст запуска субагента, привязывая живое наблюдение @@ -192,8 +214,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) { } // 2. ставим running (после валидации — чтобы плохие имена не жгли состояние) - task.Status = storage.StatusRunning - if err := w.Store.UpdateTask(ctx, task); err != nil { + if err := w.setStatus(ctx, task, storage.StatusRunning); err != nil { return fmt.Errorf("%w: set running: %v", ErrUpdate, err) } w.notifyStatus(ctx, task, storage.StatusRunning) @@ -261,16 +282,14 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) { case 0: // продолжаем на ревью case -1: - task.Status = storage.StatusTimeout - if e := w.Store.UpdateTask(ctx, task); e != nil { + if e := w.setStatus(ctx, task, storage.StatusTimeout); 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: - task.Status = storage.StatusFailed - if e := w.Store.UpdateTask(ctx, task); e != nil { + if e := w.setStatus(ctx, task, storage.StatusFailed); e != nil { return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e) } w.notifyStatus(ctx, task, storage.StatusFailed) @@ -309,9 +328,8 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) { if verdict == nil { // невалидный JSON даже после retry → failed с объяснением. - task.Status = storage.StatusFailed explain := "reviewer вернул невалидный/пустой вердикт (даже после повтора)." - if e := w.Store.UpdateTask(ctx, task); e != nil { + if e := w.setStatus(ctx, task, storage.StatusFailed); e != nil { return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e) } w.notifyStatus(ctx, task, storage.StatusFailed) @@ -325,13 +343,12 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) { 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) - } - w.notifyStatus(ctx, task, storage.StatusSuccess) - return nil + if e := w.setStatus(ctx, task, storage.StatusSuccess); 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 { @@ -341,8 +358,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) { } // Лимит исчерпан → failed с объяснением. - task.Status = storage.StatusFailed - if e := w.Store.UpdateTask(ctx, task); e != nil { + if e := w.setStatus(ctx, task, storage.StatusFailed); 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)) @@ -375,8 +391,7 @@ func (w *Worker) reviewWithRetry(ctx context.Context, taskID int64, cwd, prompt // 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 { + if e := w.setStatus(ctx, task, storage.StatusFailed); e != nil { log.Printf("worker: task %d: set failed: %v", task.ID, e) } w.notifyStatus(ctx, task, storage.StatusFailed)