feat(events): внедрить Publisher в Core/Worker/Analyst + шины в app

- 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 + логовая шина).
This commit is contained in:
ki.sagidullin
2026-08-20 04:30:29 +05:00
parent 1f7ab9a67e
commit 4957a5552e
4 changed files with 102 additions and 36 deletions

View File

@@ -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)