From 1f7ab9a67e61454e5f26dee06929ddccba48d054 Mon Sep 17 00:00:00 2001 From: "ki.sagidullin" Date: Thu, 20 Aug 2026 01:02:31 +0500 Subject: [PATCH] =?UTF-8?q?feat(events):=20EventBus=20+=20domain-=D0=BC?= =?UTF-8?q?=D0=BE=D0=B4=D0=B5=D0=BB=D1=8C=20=D1=81=D1=82=D0=B0=D1=82=D1=83?= =?UTF-8?q?=D1=81=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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). --- docs/ui-spec.md | 160 ++++++++++++++++++++++++++++++++++ internal/events/bus.go | 110 +++++++++++++++++++++++ internal/events/bus_test.go | 119 +++++++++++++++++++++++++ internal/events/events.go | 56 ++++++++++++ internal/events/log.go | 65 ++++++++++++++ internal/events/log_test.go | 76 ++++++++++++++++ internal/events/publisher.go | 20 +++++ internal/model/status.go | 77 ++++++++++++++++ internal/model/status_test.go | 50 +++++++++++ internal/storage/models.go | 85 +++++++----------- 10 files changed, 764 insertions(+), 54 deletions(-) create mode 100644 docs/ui-spec.md create mode 100644 internal/events/bus.go create mode 100644 internal/events/bus_test.go create mode 100644 internal/events/events.go create mode 100644 internal/events/log.go create mode 100644 internal/events/log_test.go create mode 100644 internal/events/publisher.go create mode 100644 internal/model/status.go create mode 100644 internal/model/status_test.go diff --git a/docs/ui-spec.md b/docs/ui-spec.md new file mode 100644 index 0000000..f69dc1a --- /dev/null +++ b/docs/ui-spec.md @@ -0,0 +1,160 @@ +# UI-спека Ratatoskr (Fyne Desktop) + +Статус: утверждается. Черновик, обсуждаем дальше (архитектура обмена UI <-> Core). + +Стек: **Go + Fyne** (свежая стабильная версия). Платформа: **Windows** (Linux пока не поддерживаем). + +--- + +## 1. Цель и роль UI + +1. UI — **полноценный инструмент** (не только админ-панель): управляет всем, чем управляет Telegram, и становится **основным** интерфейсом. +2. Один пользователь, один человек (без multi-user / мультисессии). +3. UI должен уметь **всё, что умеет Telegram-интерфейс**, и станет основным; Telegram остаётся как один из каналов. +4. UI — **графическая оболочка самого приложения** (одна программа), не отдельный process. +5. Внутри оболочки запускается **ботик сам**; UI — это лишь ещё одна реализация `chat.Channel`. +6. Запуск: **обычный запуск бинаря** → UI. Режим без UI — только по флагу **`--noui`**. + - Логика: UI включён по умолчанию, headless — опционально (`--noui`). +7. Трей-иконка не нужна. +8. При закрытии окна — **сворачивать** (не завершать процесс). Для полного выхода — **отдельная кнопка** «Завершить», чтобы бот остановился. +9. Конфиг-секции редактировать в UI пока не нужно; **hot-reload конфига** — желательно, если реализуемо (запланировать как «если возможно — да»). + +## 2. Компоновка (layout) + +- Базовый grid по умолчанию: + - **слева — список задач** (широкая колонка), + - **по центру/право — рабочая область** (в ней живёт «работа бота» / live-шаги); + - рабочая область может **разделяться на сплит-панели** (2×2, не считая левой колонки с задачами). +- Табы пока **не делаем** (первый этап). +- Панели должны быть **сворачиваемыми** (хотя бы левая колонка с задачами). +- Сохранение layout (положение/размеры панелей) — **нужно**, переживает перезапуск приложения (Fyne Preferences). +- Окно: **на весь экран по умолчанию**. +- Нижняя панель: **вкладки «Логи», «Состояние»**: + - **Логи** — ведение логов (поток системных логов); + - **Состояние** — живое состояние текущих задач / на каком статусе сейчас задача. +- Верхняя панель: **меню**, пока только пункт «О программе» (About → About). +- Тема: **тёмная**, без переключателя темы; стартуем с **дефолтной темы Fyne**. +- Локаль: **en** (интерфейс на английском). + +## 3. Данные задач + +- Задачи отображаются **списком** слева: **краткое название + служебное** (воркшоп). + - Детализация списка (статус/фильтры/поиск/сортировка/по умолчание/пагинация) — **дорабатываем потом**. +- В деталях задачи (рабочая область) показываем все перечисленные поля: + - Title, Goal, Why, AC (acceptance criteria), Repo(s), Tag (task_tag), статус, CreatedAt/UpdatedAt, вложенные трассы (traces), история чата (history), включая (session id opencode). +- История чата из БД (user/assistant) — лента сообщений. +- Трейсы агентов (аналист/dev/reviewer) — дерево шагов/ходов + raw output. +- Live-шаги (LiveRegistry) — восстановить и показывать в реальном времени. + +## 4. Действия с задачами + +- UI поддерживает все действия, доступные в TG: создание задачи, изменение, **approve**, **rework доработка**, **cancel**, **retry**, закрытие (skip/continue), статусы — в рамках валидных переходов. +- Создание задачи — форма с полями (Title, Goal/Why, Repos, AC) (какие обязательные — уточнить). +- Работа с статусами: эмуляция команд `/approve`, `/retry`, `/cancel` и т.п. (валидные переходы по state machine) +- Подтверждения опасных действий (удаление задачи, и т.п.). +- Редактирование полей (repo, AC, tags) из UI — по необходимости. + +## 7. Мониторинг системы + +- Показывать **логи** (готовый поток из `log`), **фреймы состояния** бота: + - статус opencode serve/pool, коннекции, активные сессии. +- Live-статус задач (`/status N`) эквивалентом — во вкладке «Состояние». +- Нужен индикатор занятости агента («опенкод думает/dev writing…»). + +## 8. Настройки / конфиг + +- Настройки в UI пока **не выводим**; секреты никоим образом не показываем. +- Кнопка «Завершить» (выход приложения) — в меню. +- Перезапуск бота из UI — пока не надо. + +## 9. Живое обновление / консистентность + +- **Шина событий (event-bus)** — запланирована: UI получает события из ядра (новые задачи, изменение статусов, live-шаги, логи). Явно нужен. +- UI и Telegram — **оболочки** вокруг одного ядра; обмен через **интерфейс (абстракцию)**, не зашумляя ядро. +- Консистентность между TG и UI при параллельных изменениях — через события ядра; архитектура обмена — см. раздел 12. + +## 10. Технические + +- **Fyne**: свежая стабильная (последняя). +- **Абстракция БД**: доступ к storage через **интерфейс** (не прямиком к *storage.Storage) для тестируемости. +- Пакет: **`internal/ui`** (+ подпакеты при необходимости). +- Тесты: юнит-тесты модели, **без e2e/UI-тестов** пока. +- Использование данных: **вся история** (без ограничения по времени). + +## 11. Скоуп первого этапа + +- Минимальный UI, покрывающий **функционал Telegram-интерфейса** (все команды/действия). +- Также UI — основной интерфейс (TG остаётся каналом). + +--- + +## 12. Архитектура обмена UI ↔ Core (best practices) + +### 12.1. Слои и направление зависимостей + +``` +┌──────────────┐ команды → ┌─────────────────┐ заказы/готовое → ┌──────────────┐ +│ UI (Fyne) │ ───────────▶ │ Application / │ ─────────────────▶ │ Domain/core │ +│ «представ- │ │ UseCase слой │ события (events) │ (model) │ +│ ление» │ ◀─────────── │ │ ◀────────────────── │ │ +└──────────────┘ события ← └─────────────────┘ └──────────────┘ +``` + +- **Core** (`internal/core`.ProcessTurn — state-machine) **не трогаем**. Он — источник правды. +- **App** (`internal/app`) — use-case/композиция: политика «одна активная задача на чат», `Notify`, роутинг команд. +- **UI** — ещё одна реализация `chat.Channel` (параллельно Telegram). Уже есть идеальный «port»: интерфейс `Channel{Run, OnMessage, Send, Ask, Close}`. + +### 12.2. UI получает состояние, а не управляет Core'ом (uni-directional data flow) + +- Core **мутирует состояние** (БД, worker, аналитик). UI — только **читает снимки** и реагирует на события, обновляя своё view-model (не БД и не core). +- Действия UI = **команды** (`CreateTask`, `ApproveTask`, `RetryTask`, `CancelTask`, `ContinueTask`), которые вызывают use-case в Core. +- Это даёт: единый источник правды, простую отладку, лёгкое тестирование без UI. + +### 12.3. Обмен = событийная шина (event bus / pub-sub) + +- **Core — публикует** доменные события, UI — **подписчик**: + - `TaskCreated{id, chatID}` + - `TaskStatusChanged{id, from, to}` (ready/approved/running/success/failed/timeout…) + - `TaskHistoryAppended{id, role, content}` + - `TraceAppended{id, agent, status, output}` + - `AgentActivity{id, agent, stage}` (live-шаги, индикаторы занятости) + - `LogLine` (системные логи для вкладки «Логи») +- **Направление однонаправленное**: core/worker/analyst не знают, кто подписан; никаких прямых вызовов Fyne из core! +- Реализация малая: `Bus` с `Subscribe/Unsubscribe/Publish` + buffered channels (`n`-подписчиков или `select`); паттерн уже есть в `chat.Router` (внутренний `incoming` channel). +- Сигнатуры событий использовать **и для консистентности** TG↔UI (оба канала — простые подписчики/клиенты шины). + +### 12.4. Модель потоков Fyne (v2.6+) и `fyne.Do` + +- Fyne с v2.6.0 выполняет все события/колбэки **на одной главной goroutine**. +- Фоновые goroutine (worker, opencode polling, ход core) **не должны модифицировать UI напрямую**. +- Любое обновление виджетов из своей goroutine — через: + - `fyne.Do(func(){ widget.X = …; widget.Refresh() })` — очередь в следующий кадр; + - `fyne.DoAndWait(...)` — когда надо дождаться завершения; + - унифицировать с **binding** (`binding.String`, `binding.List`) — через `fyne.Do` не требуется, но проще явно. +- Подписчик шины ставит snapshot в очередь UI-обновления и вызывает `fyne.Do`. + +### 12.5. Чтение-модель на UI: снапшоты, не мутабельные указатели + +- UI получает **копии** записей (Task, history, traces) через абстракцию БД (см. раздел 10). +- Никаких референсов на разделяемые контейнеры core; используются лёгкие view-копии (Value-типы/DT-коды) — нет гонок bottom-up. + +### 12.6. Жизненный цикл: Core ≠ UI + +- Core (бот) живёт независимо от окна. Закрытие окна = **сворачивание** (Core продолжает). +- Полный выход — только кнопка «Завершить» → `app.Quit()` (Worker.Stop → router.Close → pool.Close → store.Close). +- `--noui` → core запускается без создания окна. + +### 12.7. Скоуп реализации (поэтапно) + +- **Фаза 1**: `internal/ui` со сплитами + `TaskStore`-абстракция + **односторонний** поток: команды → Core, периодические snapshots из БД (без шины), всё через `fyne.Do`. +- **Фаза 2**: добавить `event.Bus` — Core публикует события, UI подписывается (статусы, история, live-шаги). +- **Фаза 3**: консистентность TG↔UI через общий маршрутизатор/шину (`UserID`, `chat.Router`). + +--- + +## Открытые пункты (TODO) + +- [ ] Точно определить set сплит-панелей (2x2 центр) и как добавляются +- [ ] Обязательные поля формы создания задачи +- [ ] Подробности виджета списка задач (после первого мильстон) +- [x] Архитектура обмена UI<->Core (раздел 12 — принципы; детальные контракты событий/интерфейсов — следующая итерация) \ No newline at end of file diff --git a/internal/events/bus.go b/internal/events/bus.go new file mode 100644 index 0000000..3d61901 --- /dev/null +++ b/internal/events/bus.go @@ -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) +} \ No newline at end of file diff --git a/internal/events/bus_test.go b/internal/events/bus_test.go new file mode 100644 index 0000000..906f2e5 --- /dev/null +++ b/internal/events/bus_test.go @@ -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 +} \ No newline at end of file diff --git a/internal/events/events.go b/internal/events/events.go new file mode 100644 index 0000000..88b9bd5 --- /dev/null +++ b/internal/events/events.go @@ -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() {} \ No newline at end of file diff --git a/internal/events/log.go b/internal/events/log.go new file mode 100644 index 0000000..affca1b --- /dev/null +++ b/internal/events/log.go @@ -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 +} \ No newline at end of file diff --git a/internal/events/log_test.go b/internal/events/log_test.go new file mode 100644 index 0000000..8300e06 --- /dev/null +++ b/internal/events/log_test.go @@ -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) + } +} \ No newline at end of file diff --git a/internal/events/publisher.go b/internal/events/publisher.go new file mode 100644 index 0000000..edfbaf4 --- /dev/null +++ b/internal/events/publisher.go @@ -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{} \ No newline at end of file diff --git a/internal/model/status.go b/internal/model/status.go new file mode 100644 index 0000000..ea1bef5 --- /dev/null +++ b/internal/model/status.go @@ -0,0 +1,77 @@ +// Package model — доменные типы и контракты Ratatoskr. +// +// Нижний слой архитектуры: не зависит ни от БД (storage), ни от UI (events). +// Хранит семантику предметной области (статусы задач, машина переходов), +// чтобы storage и events зависели только от этой модели, а не друг от друга. +package model + +// Status — статус задачи (state machine). +type Status string + +const ( + StatusDraft Status = "draft" // только что создана + StatusCollecting Status = "collecting" // аналитик собирает детали + StatusReady Status = "ready" // черновик готов, ждёт одобрения пользователя + StatusApproved Status = "approved" // пользователь одобрил («создавай») — воркер берёт в работу + StatusRunning Status = "running" // opencode работает + StatusSuccess Status = "success" // задача выполнена + StatusFailed Status = "failed" // ошибка выполнения + StatusTimeout Status = "timeout" // таймаут opencode + StatusCancelled Status = "cancelled" // отменена пользователем + StatusAborted Status = "aborted" // сбой сбора, черновик выброшен + StatusClosed Status = "closed" // закрыта вручную +) + +// AllStatuses — все возможные статусы для валидации. +var AllStatuses = []Status{ + StatusDraft, StatusCollecting, StatusReady, StatusApproved, + StatusRunning, StatusSuccess, StatusFailed, StatusTimeout, + StatusCancelled, StatusAborted, StatusClosed, +} + +// validTransitions задаёт разрешённые переходы статусов. +var validTransitions = map[Status][]Status{ + StatusDraft: {StatusCollecting, StatusCancelled, StatusAborted}, + StatusCollecting: {StatusReady, StatusDraft, StatusCancelled, StatusAborted}, + // ready — черновик готов: «создавай» → approved, либо правка/отмена/закрытие. + StatusReady: {StatusApproved, StatusCancelled, StatusAborted, StatusClosed, StatusCollecting}, + // approved — финальное одобрение: воркер берёт в running, либо отмена/сбой/закрытие. + StatusApproved: {StatusRunning, StatusCancelled, StatusAborted, StatusClosed}, + StatusRunning: {StatusSuccess, StatusFailed, StatusTimeout, StatusCancelled}, + StatusSuccess: {StatusClosed}, + StatusFailed: {StatusReady, StatusClosed, StatusCancelled, StatusCollecting}, // retry: перезапуск сбора + StatusTimeout: {StatusReady, StatusClosed, StatusCancelled, StatusCollecting}, // retry: перезапуск сбора + StatusCancelled: {StatusClosed}, + StatusAborted: {StatusClosed}, + StatusClosed: {}, // терминальный +} + +// IsValidTransition проверяет, допустим ли переход from → to. +func IsValidTransition(from, to Status) bool { + allowed, ok := validTransitions[from] + if !ok { + return false + } + for _, s := range allowed { + if s == to { + return true + } + } + return false +} + +// IsTerminal возвращает true, если статус терминальный. +func IsTerminal(s Status) bool { + return s == StatusSuccess || s == StatusCancelled || + s == StatusAborted || s == StatusClosed +} + +// TraceStatus — статус трассировки субагента. +type TraceStatus string + +const ( + TraceRunning TraceStatus = "running" + TraceSuccess TraceStatus = "success" + TraceFailed TraceStatus = "failed" + TraceTimeout TraceStatus = "timeout" +) \ No newline at end of file diff --git a/internal/model/status_test.go b/internal/model/status_test.go new file mode 100644 index 0000000..c824492 --- /dev/null +++ b/internal/model/status_test.go @@ -0,0 +1,50 @@ +package model + +import "testing" + +func TestValidTransitions(t *testing.T) { + cases := []struct { + from, to Status + want bool + }{ + {StatusDraft, StatusCollecting, true}, + {StatusDraft, StatusRunning, false}, + {StatusReady, StatusApproved, true}, + {StatusApproved, StatusRunning, true}, + {StatusRunning, StatusSuccess, true}, + {StatusRunning, StatusClosed, false}, + {StatusSuccess, StatusClosed, true}, + {StatusClosed, StatusDraft, false}, + } + + for _, c := range cases { + if got := IsValidTransition(c.from, c.to); got != c.want { + t.Errorf("IsValidTransition(%q, %q) = %v, want %v", c.from, c.to, got, c.want) + } + } +} + +func TestIsTerminal(t *testing.T) { + terminal := []Status{StatusSuccess, StatusCancelled, StatusAborted, StatusClosed} + for _, s := range terminal { + if !IsTerminal(s) { + t.Errorf("IsTerminal(%q) = false, want true", s) + } + } + nonTerminal := []Status{StatusDraft, StatusReady, StatusRunning, StatusCollecting, StatusApproved, StatusFailed, StatusTimeout} + for _, s := range nonTerminal { + if IsTerminal(s) { + t.Errorf("IsTerminal(%q) = true, want false", s) + } + } +} + +func TestAllStatusesCoverage(t *testing.T) { + seen := make(map[Status]bool) + for _, s := range AllStatuses { + seen[s] = true + } + if len(seen) != 11 { + t.Fatalf("AllStatuses has %d unique statuses, want 11", len(seen)) + } +} \ No newline at end of file diff --git a/internal/storage/models.go b/internal/storage/models.go index c1fe3a0..8bed32c 100644 --- a/internal/storage/models.go +++ b/internal/storage/models.go @@ -1,66 +1,43 @@ package storage -import "encoding/json" +import ( + "encoding/json" -// Status — статус задачи (state machine). -type Status string + "github.com/kamelion/ratatoskr-go/internal/model" +) +// Status — алиас доменного статуса задачи. +// +// Совместимый мост: весь внешний код продолжает использовать storage.Status +// (например «storage.StatusRunning»), но единый источник истины — model.Status. +type Status = model.Status + +// Статусы задачи — re-export из model. const ( - StatusDraft Status = "draft" // только что создана - StatusCollecting Status = "collecting" // аналитик собирает детали - StatusReady Status = "ready" // черновик готов, ждёт одобрения пользователя - StatusApproved Status = "approved" // пользователь одобрил («создавай») — воркер берёт в работу - StatusRunning Status = "running" // opencode работает - StatusSuccess Status = "success" // задача выполнена - StatusFailed Status = "failed" // ошибка выполнения - StatusTimeout Status = "timeout" // таймаут opencode - StatusCancelled Status = "cancelled" // отменена пользователем - StatusAborted Status = "aborted" // сбой сбора, черновик выброшен - StatusClosed Status = "closed" // закрыта вручную + StatusDraft = model.StatusDraft + StatusCollecting = model.StatusCollecting + StatusReady = model.StatusReady + StatusApproved = model.StatusApproved + StatusRunning = model.StatusRunning + StatusSuccess = model.StatusSuccess + StatusFailed = model.StatusFailed + StatusTimeout = model.StatusTimeout + StatusCancelled = model.StatusCancelled + StatusAborted = model.StatusAborted + StatusClosed = model.StatusClosed ) // AllStatuses — все возможные статусы для валидации. -var AllStatuses = []Status{ - StatusDraft, StatusCollecting, StatusReady, StatusApproved, - StatusRunning, StatusSuccess, StatusFailed, StatusTimeout, - StatusCancelled, StatusAborted, StatusClosed, -} - -// validTransitions задаёт разрешённые переходы статусов. -var validTransitions = map[Status][]Status{ - StatusDraft: {StatusCollecting, StatusCancelled, StatusAborted}, - StatusCollecting: {StatusReady, StatusDraft, StatusCancelled, StatusAborted}, - // ready — черновик готов: «создавай» → approved, либо правка/отмена/закрытие. - StatusReady: {StatusApproved, StatusCancelled, StatusAborted, StatusClosed, StatusCollecting}, - // approved — финальное одобрение: воркер берёт в running, либо отмена/сбой/закрытие. - StatusApproved: {StatusRunning, StatusCancelled, StatusAborted, StatusClosed}, - StatusRunning: {StatusSuccess, StatusFailed, StatusTimeout, StatusCancelled}, - StatusSuccess: {StatusClosed}, - StatusFailed: {StatusReady, StatusClosed, StatusCancelled, StatusCollecting}, // retry: перезапуск сбора - StatusTimeout: {StatusReady, StatusClosed, StatusCancelled, StatusCollecting}, // retry: перезапуск сбора - StatusCancelled: {StatusClosed}, - StatusAborted: {StatusClosed}, - StatusClosed: {}, // терминальный -} +var AllStatuses = model.AllStatuses // IsValidTransition проверяет, допустим ли переход from → to. func IsValidTransition(from, to Status) bool { - allowed, ok := validTransitions[from] - if !ok { - return false - } - for _, s := range allowed { - if s == to { - return true - } - } - return false + return model.IsValidTransition(from, to) } // IsTerminal возвращает true, если статус терминальный. func IsTerminal(s Status) bool { - return s == StatusSuccess || s == StatusCancelled || - s == StatusAborted || s == StatusClosed + return model.IsTerminal(s) } // Task — запись задачи в БД. @@ -111,14 +88,14 @@ func (t *Task) SetReposFromDB(repos string) { _ = json.Unmarshal([]byte(repos), &t.Repos) } -// Trace — запись трассировки выполнения. -type TraceStatus string +// TraceStatus — алиас доменного статуса трассировки. +type TraceStatus = model.TraceStatus const ( - TraceRunning TraceStatus = "running" - TraceSuccess TraceStatus = "success" - TraceFailed TraceStatus = "failed" - TraceTimeout TraceStatus = "timeout" + TraceRunning TraceStatus = model.TraceRunning + TraceSuccess TraceStatus = model.TraceSuccess + TraceFailed TraceStatus = model.TraceFailed + TraceTimeout TraceStatus = model.TraceTimeout ) // Trace — лог одного субагента.