feat(events): EventBus + domain-модель статусов

- internal/model: единый источник статусов задачи (Status, TraceStatus,
  машина переходов), без зависимости от storage.
- internal/storage: совместимый мост (type Status = model.Status,
  re-export констант) — внешний код не меняется.
- internal/events: шина событий (fan-out, блокирующий Publish с гарантией
  порядка), события задач/трейсов, отдельная логовая шина + LogWriter,
  Publisher/NilPublisher для внедрения в Core/Worker.
- docs/ui-spec.md: спецификация десктопного UI (Fyne).
This commit is contained in:
ki.sagidullin
2026-08-20 01:02:31 +05:00
parent 408d137747
commit 1f7ab9a67e
10 changed files with 764 additions and 54 deletions

110
internal/events/bus.go Normal file
View File

@@ -0,0 +1,110 @@
// Package events — шина событий для обмена UI ↔ Core.
//
// Две независимые шины: доменная (статусы задач, история, трейсы) и логовая
// (сырые строки лога для панели «Логи»). Обе построены на одном типе *Bus,
// DOMAIN шина блокирующая (гарантия доставки и порядка), логовая — та же,
// но с большим буфером, чтобы не тормозить логирование.
package events
import (
"sync"
)
// Event — доменное событие. Закрытый интерфейс: новые типы добавляются
// только внутри пакета.
type Event interface {
_event()
}
// Bus — широковещательная шина событий (fan-out).
//
// Publish блокирует вызывающую горутину до тех пор, пока все подписчики не
// получат событие (в копию их буфера). Порядок событий для каждого
// подписчика сохраняется. Удаление подписчика происходит горутиной-монтируется
// close(done), что снимает блокировку Publish.
type Bus struct {
bufSize int
mu sync.RWMutex
subs map[*subscriber]struct{}
}
type subscriber struct {
ch chan Event
done chan struct{}
}
// New создаёт шину с буфером bufSize на каждого подписчика.
func New(bufSize int) *Bus {
if bufSize < 1 {
bufSize = 1
}
return &Bus{
bufSize: bufSize,
subs: make(map[*subscriber]struct{}),
}
}
// Subscribe регистрирует нового подписчика и возвращает канал событий вместе
// с функцией отписки. Рекомендуемый паттерн потребления:
//
// ch, unsub := bus.Subscribe()
// defer unsub()
// for {
// select {
// case e := <-ch:
// switch ev := e.(type) { ... }
// case <-closeCh:
// return
// }
// }
//
// Канал не закрывается шиной: выход из горутины подписчика организуется через
// закрытие канала приложения либо другого сигнала в select.
func (b *Bus) Subscribe() (<-chan Event, func()) {
s := &subscriber{
ch: make(chan Event, b.bufSize),
done: make(chan struct{}),
}
b.mu.Lock()
b.subs[s] = struct{}{}
b.mu.Unlock()
var once sync.Once
unsubscribe := func() {
once.Do(func() {
close(s.done)
b.mu.Lock()
delete(b.subs, s)
b.mu.Unlock()
})
}
return s.ch, unsubscribe
}
// Publish рассылает событие всем подписчикам и блокируется, пока каждый
// подписчик либо примет событие (в свой буфер), либо отпишется. Порядок
// рассылки стабилен (по списку подписок). Безопасен для Concurrent вызовов.
func (b *Bus) Publish(e Event) {
b.mu.RLock()
subs := make([]*subscriber, 0, len(b.subs))
for s := range b.subs {
subs = append(subs, s)
}
b.mu.RUnlock()
for _, s := range subs {
select {
case s.ch <- e:
case <-s.done:
}
}
}
// SubscribersCount — число активных подписчиков (для юнит-тестов и отладки).
func (b *Bus) SubscribersCount() int {
b.mu.RLock()
defer b.mu.RUnlock()
return len(b.subs)
}

119
internal/events/bus_test.go Normal file
View File

@@ -0,0 +1,119 @@
package events
import (
"testing"
"time"
"github.com/kamelion/ratatoskr-go/internal/model"
)
// receiveOne помогает получить одно событие с таймаутом.
func receiveOne(t *testing.T, ch <-chan Event) Event {
t.Helper()
select {
case e := <-ch:
return e
case <-time.After(2 * time.Second):
t.Fatal("timeout waiting for event")
return nil
}
}
func TestBusPublishToSubscriber(t *testing.T) {
bus := New(10)
ch, unsub := bus.Subscribe()
defer unsub()
want := TaskStatusChanged{ID: 42, From: model.StatusReady, To: model.StatusRunning}
bus.Publish(want)
got := receiveOne(t, ch)
ev, ok := got.(TaskStatusChanged)
if !ok {
t.Fatalf("got %T, want TaskStatusChanged", got)
}
if ev.ID != 42 || ev.From != model.StatusReady || ev.To != model.StatusRunning {
t.Fatalf("unexpected event: %+v", ev)
}
}
func TestBusFanout(t *testing.T) {
bus := New(10)
ch1, unsub1 := bus.Subscribe()
defer unsub1()
ch2, unsub2 := bus.Subscribe()
defer unsub2()
e := LogLine{Level: "log", Text: "hello"}
bus.Publish(e)
if got := receiveOne(t, ch1); got != e {
t.Fatalf("subscriber 1 got %#v, want %#v", got, e)
}
if got := receiveOne(t, ch2); got != e {
t.Fatalf("subscriber 2 got %#v, want %#v", got, e)
}
}
func TestBusPreservesOrder(t *testing.T) {
bus := New(64)
ch, unsub := bus.Subscribe()
defer unsub()
const n = 25
for i := 0; i < n; i++ {
bus.Publish(TraceAppended{TaskID: int64(i)})
}
for i := 0; i < n; i++ {
ev := receiveOne(t, ch)
ta, ok := ev.(TraceAppended)
if !ok {
t.Fatalf("got %T, want TraceAppended", ev)
}
if ta.TaskID != int64(i) {
t.Fatalf("out of order: got %d, want %d", ta.TaskID, i)
}
}
}
func TestUnsubscribeStopsDelivery(t *testing.T) {
bus := New(10)
ch, unsub := bus.Subscribe()
bus.Publish(TaskCreated{ID: 1})
receiveOne(t, ch)
unsub()
if got := bus.SubscribersCount(); got != 0 {
t.Fatalf("SubscribersCount = %d, want 0", got)
}
// Убеждаемся, что Publish не блокируется навечно отписанным подписчиком.
bus.Publish(TaskCreated{ID: 2})
select {
case got := <-ch:
t.Fatalf("received %#v after unsubscribe", got)
case <-time.After(200 * time.Millisecond):
}
}
func TestBusSubscribersCount(t *testing.T) {
bus := New(10)
if got := bus.SubscribersCount(); got != 0 {
t.Fatalf("initial count = %d, want 0", got)
}
_, unsub1 := bus.Subscribe()
_, unsub2 := bus.Subscribe()
if got := bus.SubscribersCount(); got != 2 {
t.Fatalf("count = %d, want 2", got)
}
unsub1()
unsub2()
if got := bus.SubscribersCount(); got != 0 {
t.Fatalf("after unsub count = %d, want 0", got)
}
}
func TestNilPublisherNoOp(t *testing.T) {
NilPublisher{}.Publish(TaskCreated{ID: 1}) // must not panic
}

56
internal/events/events.go Normal file
View File

@@ -0,0 +1,56 @@
package events
import "github.com/kamelion/ratatoskr-go/internal/model"
// TaskCreated — создана новая задача.
type TaskCreated struct {
ID int64
ChatID string
Title string
}
// TaskUpdated — задача изменена (поля, права, репозитории).
type TaskUpdated struct {
ID int64
}
// TaskDeleted — задача удалена.
type TaskDeleted struct {
ID int64
}
// TaskStatusChanged — статус задачи изменился (переход из From в To).
type TaskStatusChanged struct {
ID int64
From model.Status
To model.Status
}
// HistoryAppended — добавлено сообщение в историю задачи.
type HistoryAppended struct {
TaskID int64
Role string // user | assistant | реплика события
Content string
}
// TraceAppended — добавлен/обновлён трейс субагента.
type TraceAppended struct {
TaskID int64
Agent string
Status model.TraceStatus
}
// AgentActivity — смена этапа работы агента (для анимации «состояния»).
type AgentActivity struct {
TaskID int64
Agent string
Stage string
}
func (TaskCreated) _event() {}
func (TaskUpdated) _event() {}
func (TaskDeleted) _event() {}
func (TaskStatusChanged) _event() {}
func (HistoryAppended) _event() {}
func (TraceAppended) _event() {}
func (AgentActivity) _event() {}

65
internal/events/log.go Normal file
View File

@@ -0,0 +1,65 @@
package events
import (
"strings"
"sync"
)
// LogLine — строка лога для панели «Логи».
type LogLine struct {
Level string
Text string
}
func (LogLine) _event() {}
// LogBus — тип-обёртка над *Bus для логов.
//
// Логи идут отдельной шиной, чтобы большие объёмы текста не блокировали
// доменные события и наоборот.
type LogBus struct {
*Bus
}
// NewLogBus создаёт шину логов с буфером на подписчика.
func NewLogBus(bufSize int) *LogBus {
return &LogBus{Bus: New(bufSize)}
}
// LogWriter — io.Writer, который публикует каждую строку лога как LogLine.
// Предполагается использование через log.SetOutput в связке, чтобы всё
// логирование приложения попадало и в панель «Логи».
type LogWriter struct {
bus *Bus
mu sync.Mutex // защищает остаток частичной строки
buf strings.Builder
}
// NewLogWriter создаёт LogWriter, публикующий в шину логов события LogLine{Level:"log"}.
func NewLogWriter(bus *LogBus) *LogWriter {
return &LogWriter{bus: bus.Bus}
}
// Write реализует io.Writer. Данные разрезаются по переводам строки:
// каждая законченная строка публикуется отдельным событием.
func (w *LogWriter) Write(p []byte) (int, error) {
w.mu.Lock()
defer w.mu.Unlock()
w.buf.Write(p)
data := w.buf.String()
for {
idx := strings.IndexByte(data, '\n')
if idx < 0 {
break
}
line := strings.TrimSuffix(data[:idx], "\r")
data = data[idx+1:]
if line != "" {
w.bus.Publish(LogLine{Level: "log", Text: line})
}
}
w.buf.Reset()
w.buf.WriteString(data)
return len(p), nil
}

View File

@@ -0,0 +1,76 @@
package events
import (
"strings"
"testing"
"time"
)
func TestLogWriterLines(t *testing.T) {
lbus := NewLogBus(16)
ch, unsub := lbus.Subscribe()
defer unsub()
w := NewLogWriter(lbus)
if _, err := w.Write([]byte("first line\nsecond line\n")); err != nil {
t.Fatalf("Write: %v", err)
}
first := receiveOne(t, ch)
ll, ok := first.(LogLine)
if !ok {
t.Fatalf("got %T, want LogLine", first)
}
if ll.Text != "first line" || ll.Level != "log" {
t.Fatalf("unexpected first line: %+v", ll)
}
second := receiveOne(t, ch)
if ll, ok := second.(LogLine); !ok || ll.Text != "second line" {
t.Fatalf("unexpected second line: %+v", second)
}
}
func TestLogWriterPartialLine(t *testing.T) {
lbus := NewLogBus(16)
ch, unsub := lbus.Subscribe()
defer unsub()
w := NewLogWriter(lbus)
if _, err := w.Write([]byte("partial")); err != nil {
t.Fatalf("Write: %v", err)
}
select {
case got := <-ch:
t.Fatalf("partial line should not be published yet, got %#v", got)
case <-time.After(150 * time.Millisecond):
}
if _, err := w.Write([]byte(" line\n")); err != nil {
t.Fatalf("Write: %v", err)
}
got := receiveOne(t, ch)
ll, ok := got.(LogLine)
if !ok || ll.Text != "partial line" {
t.Fatalf("unexpected: %#v", got)
}
}
func TestLogWriterMultiSplit(t *testing.T) {
lbus := NewLogBus(16)
ch, unsub := lbus.Subscribe()
defer unsub()
w := NewLogWriter(lbus)
if _, err := w.Write([]byte("line1\nline2\nline3\n")); err != nil {
t.Fatalf("Write: %v", err)
}
var texts []string
for i := 0; i < 3; i++ {
ev := receiveOne(t, ch)
texts = append(texts, ev.(LogLine).Text)
}
if strings.Join(texts, ",") != "line1,line2,line3" {
t.Fatalf("got %v", texts)
}
}

View File

@@ -0,0 +1,20 @@
package events
// Publisher — минимальный интерфейс для встраивания шины в Core/Worker/Analyst.
//
// Core публикует доменные события через него; конкретная шина подставляется
// при сборке приложения. На время тестов или до создания UI можно
// использовать NilPublisher — безопасную no-op реализацию.
type Publisher interface {
Publish(e Event)
}
// NilPublisher — no-op издатель, чтобы компоненты могли работать без UI.
type NilPublisher struct{}
// Publish ничего не делает (совместимо с Publisher).
func (NilPublisher) Publish(Event) {}
// CompileTime-проверка: *Bus реализует Publisher.
var _ Publisher = (*Bus)(nil)
var _ Publisher = NilPublisher{}