From 1f7ab9a67e61454e5f26dee06929ddccba48d054 Mon Sep 17 00:00:00 2001 From: "ki.sagidullin" Date: Thu, 20 Aug 2026 01:02:31 +0500 Subject: [PATCH 1/5] =?UTF-8?q?feat(events):=20EventBus=20+=20domain-?= =?UTF-8?q?=D0=BC=D0=BE=D0=B4=D0=B5=D0=BB=D1=8C=20=D1=81=D1=82=D0=B0=D1=82?= =?UTF-8?q?=D1=83=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 — лог одного субагента. -- 2.49.1 From 4957a5552eb3524f59e1773441910d931aacac2c Mon Sep 17 00:00:00 2001 From: "ki.sagidullin" Date: Thu, 20 Aug 2026 04:30:29 +0500 Subject: [PATCH 2/5] =?UTF-8?q?feat(events):=20=D0=B2=D0=BD=D0=B5=D0=B4?= =?UTF-8?q?=D1=80=D0=B8=D1=82=D1=8C=20Publisher=20=D0=B2=20Core/Worker/Ana?= =?UTF-8?q?lyst=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) -- 2.49.1 From 66e40d384e812427f84617147793e3329e40a456 Mon Sep 17 00:00:00 2001 From: "ki.sagidullin" Date: Thu, 20 Aug 2026 07:17:57 +0500 Subject: [PATCH 3/5] =?UTF-8?q?feat(events):=20=D1=82=D0=B8=D0=BF=D0=B8?= =?UTF-8?q?=D0=B7=D0=B8=D1=80=D0=BE=D0=B2=D0=B0=D0=BD=D0=BD=D1=8B=D0=B9=20?= =?UTF-8?q?Hub=20=E2=80=94=20=D0=BF=D0=BE=D0=B4=D0=BF=D0=B8=D1=81=D1=87?= =?UTF-8?q?=D0=B8=D0=BA=20=D1=88=D0=B8=D0=BD=D1=8B=20=D1=81=20=D0=B4=D0=B8?= =?UTF-8?q?=D1=81=D0=BF=D0=B5=D1=82=D1=87=D0=B5=D1=80=D0=B8=D0=B7=D0=B0?= =?UTF-8?q?=D1=86=D0=B8=D0=B5=D0=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Hub слушает *Bus в собственной горутине и вызывает зарегистрированные обработчики по типу события (On[T]); порядок сохраняется. - Мост к UI: колбэки выполняются в горутине Hub → внутри можно переложить работу на поток Fyne (fyne.Do) или thread-safe binding. - Close останавливает горутину и отписывается; Start идемпотентен. - Юнит-тесты: доставка по типу, игнор посторонних типов, несколько обработчиков одного типа, остановка после Close. --- internal/events/hub.go | 96 +++++++++++++++++++++++++++++ internal/events/hub_test.go | 119 ++++++++++++++++++++++++++++++++++++ 2 files changed, 215 insertions(+) create mode 100644 internal/events/hub.go create mode 100644 internal/events/hub_test.go diff --git a/internal/events/hub.go b/internal/events/hub.go new file mode 100644 index 0000000..51aeaaa --- /dev/null +++ b/internal/events/hub.go @@ -0,0 +1,96 @@ +package events + +import ( + "reflect" + "sync" +) + +// Handler — колбэк-обработчик события конкретного типа. +type Handler[T Event] func(e T) + +// Hub — типизированный подписчик шины. +// +// Слушает *Bus в собственной горутине и вызывает зарегистрированные обработчики +// для событий соответствующих типов. Порядок обработки сохраняется (порядок +// шины). Обработчики одного типа вызываются в порядке регистрации. +// +// Это мост между шиной и UI-потоком: колбэки выполняются в горутине Hub, поэтому +// внутри них нужно либо перекладывать работу на поток Fyne (fyne.Do), либо +// пользоваться только thread-safe структурами (binding). +type Hub struct { + bus *Bus + done chan struct{} + once sync.Once + ch <-chan Event + unsub func() + + mu sync.Mutex + handlers map[reflect.Type][]any +} + +// NewHub создаёт подписчик на указанную шину (пока не запущен). +func NewHub(bus *Bus) *Hub { + return &Hub{ + bus: bus, + done: make(chan struct{}), + handlers: make(map[reflect.Type][]any), + } +} + +// On регистрирует обработчик для типа события T. Безопасно вызывать до Start +// и из других горутин. +func On[T Event](h *Hub, fn Handler[T]) { + h.mu.Lock() + defer h.mu.Unlock() + var zero T + typ := reflect.TypeOf(zero) + h.handlers[typ] = append(h.handlers[typ], fn) +} + +// Start подписывается на шину и запускает горутину чтения событий. +func (h *Hub) Start() { + if h.ch != nil { + return + } + h.ch, h.unsub = h.bus.Subscribe() + go h.run() +} + +// Close останавливает горутину и отписывается от шины. Идемпотентен. +func (h *Hub) Close() { + h.once.Do(func() { + close(h.done) + if h.unsub != nil { + h.unsub() + } + }) +} + +// run — цикл чтения событий и диспетчеризации. +func (h *Hub) run() { + defer h.Close() + for { + select { + case e, ok := <-h.ch: + if !ok { + return + } + h.dispatch(e) + case <-h.done: + return + } + } +} + +// dispatch вызывает все обработчики, зарегистрированные для типа события e. +func (h *Hub) dispatch(e Event) { + typ := reflect.TypeOf(e) + + h.mu.Lock() + fns := append([]any(nil), h.handlers[typ]...) + h.mu.Unlock() + + for _, fn := range fns { + reflect.ValueOf(fn).Call([]reflect.Value{reflect.ValueOf(e)}) + } +} \ No newline at end of file diff --git a/internal/events/hub_test.go b/internal/events/hub_test.go new file mode 100644 index 0000000..04abb1c --- /dev/null +++ b/internal/events/hub_test.go @@ -0,0 +1,119 @@ +package events + +import ( + "sync/atomic" + "testing" + "time" + + "github.com/kamelion/ratatoskr-go/internal/model" +) + +func TestHubDeliversTypedEvent(t *testing.T) { + bus := New(16) + h := NewHub(bus) + + var got atomic.Value + On[TaskStatusChanged](h, func(e TaskStatusChanged) { + got.CompareAndSwap(nil, e) + }) + + h.Start() + defer h.Close() + + want := TaskStatusChanged{ID: 7, From: model.StatusReady, To: model.StatusRunning} + bus.Publish(want) + + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + if raw := got.Load(); raw != nil { + ev := raw.(TaskStatusChanged) + if ev.ID != 7 || ev.From != model.StatusReady || ev.To != model.StatusRunning { + t.Fatalf("unexpected event: %+v", ev) + } + return + } + time.Sleep(10 * time.Millisecond) + } + t.Fatal("handler was not called") +} + +func TestHubIgnoresUnrelatedTypes(t *testing.T) { + bus := New(16) + h := NewHub(bus) + + var calls atomic.Int32 + On[TaskStatusChanged](h, func(TaskStatusChanged) { calls.Add(1) }) + + h.Start() + defer h.Close() + + // другие типы событий не должны дойти до этого обработчика + bus.Publish(AgentActivity{TaskID: 1, Agent: "dev", Stage: "run"}) + bus.Publish(LogLine{Level: "log", Text: "x"}) + + time.Sleep(200 * time.Millisecond) + if n := calls.Load(); n != 0 { + t.Fatalf("handler called %d times for unrelated events", n) + } +} + +func TestHubMultipleHandlersSameType(t *testing.T) { + bus := New(16) + h := NewHub(bus) + + var a, b atomic.Int32 + On[HistoryAppended](h, func(HistoryAppended) { a.Add(1) }) + On[HistoryAppended](h, func(HistoryAppended) { b.Add(1) }) + + h.Start() + defer h.Close() + + bus.Publish(HistoryAppended{TaskID: 1, Role: "user", Content: "hi"}) + + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + if a.Load() == 1 && b.Load() == 1 { + return + } + time.Sleep(10 * time.Millisecond) + } + t.Fatal("not all handlers called") +} + +func TestHubCloseStopsDelivery(t *testing.T) { + bus := New(16) + h := NewHub(bus) + + var calls atomic.Int32 + On[TaskCreated](h, func(TaskCreated) { calls.Add(1) }) + + h.Start() + bus.Publish(TaskCreated{ID: 1}) + + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) && calls.Load() == 0 { + time.Sleep(10 * time.Millisecond) + } + if calls.Load() == 0 { + t.Fatal("initial delivery failed") + } + + h.Close() + // после Close Hub отписан — события не доходят + bus.Publish(TaskCreated{ID: 2}) + time.Sleep(150 * time.Millisecond) + if n := calls.Load(); n > 1 { + t.Fatalf("handler called %d times after Close", n) + } +} + +func TestHubStartIdempotent(t *testing.T) { + bus := New(16) + h := NewHub(bus) + h.Start() + h.Start() // повторный Start не должен создавать вторую горутину + h.Close() + // если бы было две горутины — Publish блокировался бы на буфере, но закрытие сняло бы блок; + // просто проверяем, что Close и повторный Start не падают + bus.Publish(TaskCreated{ID: 1}) +} \ No newline at end of file -- 2.49.1 From baf2ad7147589be56fd2c8e93f23429b95b5b2dd Mon Sep 17 00:00:00 2001 From: "ki.sagidullin" Date: Thu, 20 Aug 2026 08:15:42 +0500 Subject: [PATCH 4/5] =?UTF-8?q?feat(ui):=20=D0=BA=D0=BE=D0=BD=D1=82=D1=80?= =?UTF-8?q?=D0=BE=D0=BB=D0=BB=D0=B5=D1=80=20=E2=80=94=20=D0=BF=D0=BE=D0=B4?= =?UTF-8?q?=D0=BF=D0=B8=D1=81=D1=87=D0=B8=D0=BA=20=D1=88=D0=B8=D0=BD=20?= =?UTF-8?q?=D0=B8=20=D0=BC=D0=BE=D1=81=D1=82=20=D0=BA=20Fyne-=D0=BF=D1=80?= =?UTF-8?q?=D0=B5=D0=B4=D1=81=D1=82=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD=D0=B8?= =?UTF-8?q?=D1=8E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - View — интерфейс Fyne-слоя: колбэки на все доменные события + логи. - Controller подписывается на доменную и логовую шины (через events.Hub) и диспатчит события в View; колбэки выполняются в горутинах Hub, поэтому реализация Fyne обязана обновлять виджеты через fyne.Do (спец 12.4). - Close отписывает контроллер (сворачивание окна) без остановки Core. - NilView — no-op для headless (--noui) и тестов. - Юнит-тесты: доставка доменных/логовых событий, остановка после Close. --- internal/ui/controller.go | 52 ++++++++++++ internal/ui/controller_test.go | 146 +++++++++++++++++++++++++++++++++ internal/ui/view.go | 40 +++++++++ 3 files changed, 238 insertions(+) create mode 100644 internal/ui/controller.go create mode 100644 internal/ui/controller_test.go create mode 100644 internal/ui/view.go diff --git a/internal/ui/controller.go b/internal/ui/controller.go new file mode 100644 index 0000000..b53f606 --- /dev/null +++ b/internal/ui/controller.go @@ -0,0 +1,52 @@ +package ui + +import ( + "github.com/kamelion/ratatoskr-go/internal/events" +) + +// Controller — подписчик шин событий, мост к Fyne-представлению. +// +// Регистрирует Hub'ы на доменной и логовой шинах и диспатчит события в View +// (см. спец 12.4): колбэки View выполняются в горутинах Hub, поэтому +// реализация Fyne должна вызывать fyne.Do, чтобы обновить виджеты на главной +// горутине окна. +// +// Жизненный цикл независим от ядра: Close только отписывает Controller от шин +// (например, при сворачивании окна), но Core/worker продолжают работать. +type Controller struct { + domain *events.Hub + logs *events.Hub + view View +} + +// New создаёт Controller, подписанный на доменную (domainBus) и логовую +// (logBus) шины, с представлением view. +func New(domainBus *events.Bus, logBus *events.LogBus, view View) *Controller { + c := &Controller{view: view} + + c.domain = events.NewHub(domainBus) + events.On(c.domain, func(e events.TaskCreated) { view.OnTaskCreated(e) }) + events.On(c.domain, func(e events.TaskUpdated) { view.OnTaskUpdated(e) }) + events.On(c.domain, func(e events.TaskDeleted) { view.OnTaskDeleted(e) }) + events.On(c.domain, func(e events.TaskStatusChanged) { view.OnTaskStatusChanged(e) }) + events.On(c.domain, func(e events.HistoryAppended) { view.OnHistoryAppended(e) }) + events.On(c.domain, func(e events.TraceAppended) { view.OnTraceAppended(e) }) + events.On(c.domain, func(e events.AgentActivity) { view.OnAgentActivity(e) }) + + c.logs = events.NewHub(logBus.Bus) + events.On(c.logs, func(e events.LogLine) { view.OnLog(e) }) + + return c +} + +// Start запускает подписку на обе шины. +func (c *Controller) Start() { + c.domain.Start() + c.logs.Start() +} + +// Close останавливает подписки. Идемпотентен. +func (c *Controller) Close() { + c.domain.Close() + c.logs.Close() +} \ No newline at end of file diff --git a/internal/ui/controller_test.go b/internal/ui/controller_test.go new file mode 100644 index 0000000..67df3f9 --- /dev/null +++ b/internal/ui/controller_test.go @@ -0,0 +1,146 @@ +package ui + +import ( + "sync" + "testing" + "time" + + "github.com/kamelion/ratatoskr-go/internal/events" + "github.com/kamelion/ratatoskr-go/internal/model" +) + +// recordingView фиксирует события, дошедшие до View. +type recordingView struct { + mu sync.Mutex + + created int + updated int + deleted int + status []events.TaskStatusChanged + history []events.HistoryAppended + traces int + act []events.AgentActivity + logs []events.LogLine +} + +func (v *recordingView) OnTaskCreated(events.TaskCreated) { v.created++ } +func (v *recordingView) OnTaskUpdated(events.TaskUpdated) { v.updated++ } +func (v *recordingView) OnTaskDeleted(events.TaskDeleted) { v.deleted++ } +func (v *recordingView) OnHistoryAppended(e events.HistoryAppended) { + v.mu.Lock() + v.history = append(v.history, e) + v.mu.Unlock() +} +func (v *recordingView) OnTraceAppended(events.TraceAppended) { + v.mu.Lock() + v.traces++ + v.mu.Unlock() +} +func (v *recordingView) OnTaskStatusChanged(e events.TaskStatusChanged) { + v.mu.Lock() + v.status = append(v.status, e) + v.mu.Unlock() +} +func (v *recordingView) OnAgentActivity(e events.AgentActivity) { + v.mu.Lock() + v.act = append(v.act, e) + v.mu.Unlock() +} +func (v *recordingView) OnLog(e events.LogLine) { + v.mu.Lock() + v.logs = append(v.logs, e) + v.mu.Unlock() +} + +// waitFor ждёт, пока условие не станет истинным (до 2 секунд). +func waitFor(cond func() bool) bool { + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + if cond() { + return true + } + time.Sleep(10 * time.Millisecond) + } + return cond() +} + +func TestControllerDeliversDomainEvents(t *testing.T) { + dbus := events.New(64) + lbus := events.NewLogBus(64) + view := &recordingView{} + c := New(dbus, lbus, view) + c.Start() + defer c.Close() + + dbus.Publish(events.TaskStatusChanged{ID: 1, From: model.StatusReady, To: model.StatusRunning}) + dbus.Publish(events.HistoryAppended{TaskID: 1, Role: "user", Content: "hello"}) + dbus.Publish(events.AgentActivity{TaskID: 1, Agent: "dev", Stage: "run"}) + + if !waitFor(func() bool { return len(view.status) == 1 }) { + t.Fatal("TaskStatusChanged not delivered") + } + if !waitFor(func() bool { return len(view.history) == 1 }) { + t.Fatal("HistoryAppended not delivered") + } + if !waitFor(func() bool { return len(view.act) == 1 }) { + t.Fatal("AgentActivity not delivered") + } + + sc := view.status[0] + if sc.ID != 1 || sc.From != model.StatusReady || sc.To != model.StatusRunning { + t.Fatalf("bad status event: %+v", sc) + } +} + +func TestControllerDeliversLogEvents(t *testing.T) { + dbus := events.New(64) + lbus := events.NewLogBus(64) + view := &recordingView{} + c := New(dbus, lbus, view) + c.Start() + defer c.Close() + + lbus.Bus.Publish(events.LogLine{Level: "log", Text: "worker started"}) + + if !waitFor(func() bool { return len(view.logs) == 1 }) { + t.Fatal("LogLine not delivered") + } + if got := view.logs[0]; got.Text != "worker started" || got.Level != "log" { + t.Fatalf("bad log event: %+v", got) + } +} + +func TestControllerCloseStopsDelivery(t *testing.T) { + dbus := events.New(64) + lbus := events.NewLogBus(64) + view := &recordingView{} + c := New(dbus, lbus, view) + c.Start() + + dbus.Publish(events.TaskCreated{ID: 1}) + if !waitFor(func() bool { return view.created == 1 }) { + t.Fatal("initial delivery failed") + } + + c.Close() + + dbus.Publish(events.TaskCreated{ID: 2}) + lbus.Bus.Publish(events.LogLine{Level: "log", Text: "after close"}) + time.Sleep(150 * time.Millisecond) + + if view.created > 1 { + t.Fatalf("TaskCreated delivered %d times after Close", view.created) + } + if len(view.logs) > 0 { + t.Fatalf("LogLine delivered %d times after Close", len(view.logs)) + } +} + +func TestControllerIdempotentClose(t *testing.T) { + dbus := events.New(16) + lbus := events.NewLogBus(16) + c := New(dbus, lbus, &recordingView{}) + c.Start() + c.Close() + c.Close() // не должно падать +} \ No newline at end of file diff --git a/internal/ui/view.go b/internal/ui/view.go new file mode 100644 index 0000000..38df43c --- /dev/null +++ b/internal/ui/view.go @@ -0,0 +1,40 @@ +// Package ui — контроллер и представления десктопного интерфейса (Fyne). +// +// Слой обмена с ядром: однонаправленный поток (спец 12.2–12.5). +// - Core мутирует состояние; UI только читает снимки и реагирует на события. +// - Действия UI = команды (CreateTask/ApproveTask/...), которые зовут Core. +// - Подписчик шины (Controller) получает события и перекладывает их в View +// (реализация Fyne) через fyne.Do — никаких прямых вызовов Fyne из core. +package ui + +import ( + "github.com/kamelion/ratatoskr-go/internal/events" +) + +// View — интерфейс, который реализует Fyne-слой окна. +// +// Методы вызываются из горутины Controller (горутина Hub-подписчика), поэтому +// реализация обязана перекладывать работу на поток Fyne через fyne.Do / +// fyne.DoAndWait, либо использовать thread-safe структуры (binding). +type View interface { + OnTaskCreated(e events.TaskCreated) + OnTaskUpdated(e events.TaskUpdated) + OnTaskDeleted(e events.TaskDeleted) + OnTaskStatusChanged(e events.TaskStatusChanged) + OnHistoryAppended(e events.HistoryAppended) + OnTraceAppended(e events.TraceAppended) + OnAgentActivity(e events.AgentActivity) + OnLog(e events.LogLine) +} + +// NilView — no-op реализация View для headless-режима и тестов. +type NilView struct{} + +func (NilView) OnTaskCreated(events.TaskCreated) {} +func (NilView) OnTaskUpdated(events.TaskUpdated) {} +func (NilView) OnTaskDeleted(events.TaskDeleted) {} +func (NilView) OnTaskStatusChanged(events.TaskStatusChanged) {} +func (NilView) OnHistoryAppended(events.HistoryAppended) {} +func (NilView) OnTraceAppended(events.TraceAppended) {} +func (NilView) OnAgentActivity(events.AgentActivity) {} +func (NilView) OnLog(events.LogLine) {} \ No newline at end of file -- 2.49.1 From c2272137b3cdb2a3b75f73fb60261720061aa9e4 Mon Sep 17 00:00:00 2001 From: "ki.sagidullin" Date: Thu, 20 Aug 2026 08:54:55 +0500 Subject: [PATCH 5/5] =?UTF-8?q?feat(ui):=20=D0=BA=D0=BE=D0=BC=D0=B0=D0=BD?= =?UTF-8?q?=D0=B4=D1=8B,=20snapshots,=20Fyne-=D0=BE=D0=BA=D0=BD=D0=BE=20?= =?UTF-8?q?=D0=B8=20=D0=B8=D0=BD=D1=82=D0=B5=D0=B3=D1=80=D0=B0=D1=86=D0=B8?= =?UTF-8?q?=D1=8F=20(--noui)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Команды (спец 12.2): - ui.Commands: действия UI → текстовые команды канала (start/cancel/skip/ retry/continue/approve/send), единый путь через chat.Channel. - ui.Window = chat.Channel + View; NilWindow для headless/тестов. Snapshots (спец 12.5): - ui.Store: чтение-модель, возвращает только копии (ListTasks/GetTask/ GetHistory/GetTraces); ui.DBStore поверх storage. Fyne-окно (internal/ui/desktop, build-tag cgo): - список задач слева, сплиты рабочей области и панели «Логи»/«Состояние», ввод+кнопки команд, тёмная тема, fullscreen, сохранение layout в Preferences, сворачивание при закрытии крестиком. - Колбэки View через fyne.Do (спец 12.4). Интеграция: - app.New(..., noUI); флаг --noui; UI собирается только с cgo (ui_cgo/ ui_noui фабрики), приложение headless без него. - Run: окно блокирует главную горутину; «Завершить» → cancel → graceful shutdown (спец 12.6). - go.mod: fyne.io/fyne/v2 v2.6.0 (direct). --- cmd/ratatoskr/main.go | 5 +- docs/ui-spec.md | 13 +- go.mod | 31 ++ go.sum | 74 ++++- internal/app/app.go | 32 ++- internal/app/app_test.go | 16 +- internal/app/ui_cgo.go | 15 + internal/app/ui_noui.go | 11 + internal/ui/commands.go | 46 +++ internal/ui/controller.go | 2 +- internal/ui/controller_test.go | 2 +- internal/ui/desktop/window.go | 497 +++++++++++++++++++++++++++++++++ internal/ui/store.go | 73 +++++ internal/ui/store_test.go | 107 +++++++ internal/ui/view.go | 2 +- internal/ui/window.go | 68 +++++ internal/ui/window_test.go | 80 ++++++ 17 files changed, 1056 insertions(+), 18 deletions(-) create mode 100644 internal/app/ui_cgo.go create mode 100644 internal/app/ui_noui.go create mode 100644 internal/ui/commands.go create mode 100644 internal/ui/desktop/window.go create mode 100644 internal/ui/store.go create mode 100644 internal/ui/store_test.go create mode 100644 internal/ui/window.go create mode 100644 internal/ui/window_test.go diff --git a/cmd/ratatoskr/main.go b/cmd/ratatoskr/main.go index ea1d33c..7e549f9 100644 --- a/cmd/ratatoskr/main.go +++ b/cmd/ratatoskr/main.go @@ -27,6 +27,7 @@ var updateToken = "" func main() { cfg := flag.String("config", "", "путь к config.yaml (по умолчанию — CWD/config.yaml)") versionFlag := flag.Bool("version", false, "показать версию и выйти") + noUI := flag.Bool("noui", false, "работать без графического окна (headless)") flag.Parse() if *versionFlag { @@ -34,7 +35,7 @@ func main() { return } - a, err := app.New(*cfg, version, updateToken) + a, err := app.New(*cfg, version, updateToken, *noUI) if err != nil { log.Fatalf("app init: %v", err) } @@ -60,4 +61,4 @@ func main() { } log.Print("ratatoskr: остановлен") os.Exit(0) -} \ No newline at end of file +} diff --git a/docs/ui-spec.md b/docs/ui-spec.md index f69dc1a..363fb09 100644 --- a/docs/ui-spec.md +++ b/docs/ui-spec.md @@ -146,9 +146,16 @@ ### 12.7. Скоуп реализации (поэтапно) -- **Фаза 1**: `internal/ui` со сплитами + `TaskStore`-абстракция + **односторонний** поток: команды → Core, периодические snapshots из БД (без шины), всё через `fyne.Do`. -- **Фаза 2**: добавить `event.Bus` — Core публикует события, UI подписывается (статусы, история, live-шаги). -- **Фаза 3**: консистентность TG↔UI через общий маршрутизатор/шину (`UserID`, `chat.Router`). +- [x] **Фаза 1**: `internal/ui` — Store-абстракция (копии), окно как `chat.Channel`, + команды через канал, snapshots из БД. Fyne-реализация — `internal/ui/desktop` + (build-tag `cgo`, требует компилятор C для GLFW). +- [x] **Фаза 2**: `internal/events` — Bus/Hub, Core публикует события, UI подписывается + через `Controller` (статусы, история, live-шаги, логи). +- [ ] **Фаза 3**: консистентность TG↔UI через общий маршрутизатор/шину (`UserID`, `chat.Router`). + +Примечание: окно (Фаза 1) собирается только с cgo; на машинах без C-компилятора +приложение работает headless (UI=nil). Реальное окно полноценно проверяется +на машине с MinGW-w64/gcc (`go build -ldflags ...`), здесь — только typecheck. --- diff --git a/go.mod b/go.mod index 6f02588..4bb5bbe 100644 --- a/go.mod +++ b/go.mod @@ -3,17 +3,48 @@ module github.com/kamelion/ratatoskr-go go 1.25.0 require ( + fyne.io/fyne/v2 v2.6.0 gopkg.in/yaml.v3 v3.0.1 modernc.org/sqlite v1.56.0 ) require ( + fyne.io/systray v1.11.0 // indirect + github.com/BurntSushi/toml v1.4.0 // indirect + github.com/davecgh/go-spew v1.1.1 // indirect github.com/dustin/go-humanize v1.0.1 // indirect + github.com/fredbi/uri v1.1.0 // indirect + github.com/fsnotify/fsnotify v1.7.0 // indirect + github.com/fyne-io/gl-js v0.1.0 // indirect + github.com/fyne-io/glfw-js v0.2.0 // indirect + github.com/fyne-io/image v0.1.1 // indirect + github.com/fyne-io/oksvg v0.1.0 // indirect + github.com/go-gl/gl v0.0.0-20231021071112-07e5d0ea2e71 // indirect + github.com/go-gl/glfw/v3.3/glfw v0.0.0-20240506104042-037f3cc74f2a // indirect + github.com/go-text/render v0.2.0 // indirect + github.com/go-text/typesetting v0.2.1 // indirect + github.com/godbus/dbus/v5 v5.1.0 // indirect github.com/google/uuid v1.6.0 // indirect + github.com/hack-pad/go-indexeddb v0.3.2 // indirect + github.com/hack-pad/safejs v0.1.0 // indirect + github.com/jeandeaual/go-locale v0.0.0-20241217141322-fcc2cadd6f08 // indirect + github.com/jsummers/gobmp v0.0.0-20230614200233-a9de23ed2e25 // indirect + github.com/kr/text v0.2.0 // indirect github.com/mattn/go-isatty v0.0.24 // indirect github.com/ncruces/go-strftime v1.0.0 // indirect + github.com/nfnt/resize v0.0.0-20180221191011-83c6a9932646 // indirect + github.com/nicksnyder/go-i18n/v2 v2.5.1 // indirect + github.com/pmezard/go-difflib v1.0.0 // indirect github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect + github.com/rymdport/portal v0.4.1 // indirect + github.com/srwiley/oksvg v0.0.0-20221011165216-be6e8873101c // indirect + github.com/srwiley/rasterx v0.0.0-20220730225603-2ab79fcdd4ef // indirect + github.com/stretchr/testify v1.10.0 // indirect + github.com/yuin/goldmark v1.7.8 // indirect + golang.org/x/image v0.24.0 // indirect + golang.org/x/net v0.35.0 // indirect golang.org/x/sys v0.47.0 // indirect + golang.org/x/text v0.22.0 // indirect modernc.org/libc v1.74.4 // indirect modernc.org/mathutil v1.7.1 // indirect modernc.org/memory v1.11.0 // indirect diff --git a/go.sum b/go.sum index a6bc08b..babcfb7 100644 --- a/go.sum +++ b/go.sum @@ -1,27 +1,99 @@ +fyne.io/fyne/v2 v2.6.0 h1:Rywo9yKYN4qvNuvkRuLF+zxhJYWbIFM+m4N4KV4p1pQ= +fyne.io/fyne/v2 v2.6.0/go.mod h1:YZt7SksjvrSNJCwbWFV32WON3mE1Sr7L41D29qMZ/lU= +fyne.io/systray v1.11.0 h1:D9HISlxSkx+jHSniMBR6fCFOUjk1x/OOOJLa9lJYAKg= +fyne.io/systray v1.11.0/go.mod h1:RVwqP9nYMo7h5zViCBHri2FgjXF7H2cub7MAq4NSoLs= +github.com/BurntSushi/toml v1.4.0 h1:kuoIxZQy2WRRk1pttg9asf+WVv6tWQuBNVmK8+nqPr0= +github.com/BurntSushi/toml v1.4.0/go.mod h1:ukJfTF/6rtPPRCnwkur4qwRxa8vTRFBF0uk2lLoLwho= +github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= +github.com/felixge/fgprof v0.9.3 h1:VvyZxILNuCiUCSXtPtYmmtGvb65nqXh2QFWc0Wpf2/g= +github.com/felixge/fgprof v0.9.3/go.mod h1:RdbpDgzqYVh/T9fPELJyV7EYJuHB55UTEULNun8eiPw= +github.com/fredbi/uri v1.1.0 h1:OqLpTXtyRg9ABReqvDGdJPqZUxs8cyBDOMXBbskCaB8= +github.com/fredbi/uri v1.1.0/go.mod h1:aYTUoAXBOq7BLfVJ8GnKmfcuURosB1xyHDIfWeC/iW4= +github.com/fsnotify/fsnotify v1.7.0 h1:8JEhPFa5W2WU7YfeZzPNqzMP6Lwt7L2715Ggo0nosvA= +github.com/fsnotify/fsnotify v1.7.0/go.mod h1:40Bi/Hjc2AVfZrqy+aj+yEI+/bRxZnMJyTJwOpGvigM= +github.com/fyne-io/gl-js v0.1.0 h1:8luJzNs0ntEAJo+8x8kfUOXujUlP8gB3QMOxO2mUdpM= +github.com/fyne-io/gl-js v0.1.0/go.mod h1:ZcepK8vmOYLu96JoxbCKJy2ybr+g1pTnaBDdl7c3ajI= +github.com/fyne-io/glfw-js v0.2.0 h1:8GUZtN2aCoTPNqgRDxK5+kn9OURINhBEBc7M4O1KrmM= +github.com/fyne-io/glfw-js v0.2.0/go.mod h1:Ri6te7rdZtBgBpxLW19uBpp3Dl6K9K/bRaYdJ22G8Jk= +github.com/fyne-io/image v0.1.1 h1:WH0z4H7qfvNUw5l4p3bC1q70sa5+YWVt6HCj7y4VNyA= +github.com/fyne-io/image v0.1.1/go.mod h1:xrfYBh6yspc+KjkgdZU/ifUC9sPA5Iv7WYUBzQKK7JM= +github.com/fyne-io/oksvg v0.1.0 h1:7EUKk3HV3Y2E+qypp3nWqMXD7mum0hCw2KEGhI1fnBw= +github.com/fyne-io/oksvg v0.1.0/go.mod h1:dJ9oEkPiWhnTFNCmRgEze+YNprJF7YRbpjgpWS4kzoI= +github.com/go-gl/gl v0.0.0-20231021071112-07e5d0ea2e71 h1:5BVwOaUSBTlVZowGO6VZGw2H/zl9nrd3eCZfYV+NfQA= +github.com/go-gl/gl v0.0.0-20231021071112-07e5d0ea2e71/go.mod h1:9YTyiznxEY1fVinfM7RvRcjRHbw2xLBJ3AAGIT0I4Nw= +github.com/go-gl/glfw/v3.3/glfw v0.0.0-20240506104042-037f3cc74f2a h1:vxnBhFDDT+xzxf1jTJKMKZw3H0swfWk9RpWbBbDK5+0= +github.com/go-gl/glfw/v3.3/glfw v0.0.0-20240506104042-037f3cc74f2a/go.mod h1:tQ2UAYgL5IevRw8kRxooKSPJfGvJ9fJQFa0TUsXzTg8= +github.com/go-text/render v0.2.0 h1:LBYoTmp5jYiJ4NPqDc2pz17MLmA3wHw1dZSVGcOdeAc= +github.com/go-text/render v0.2.0/go.mod h1:CkiqfukRGKJA5vZZISkjSYrcdtgKQWRa2HIzvwNN5SU= +github.com/go-text/typesetting v0.2.1 h1:x0jMOGyO3d1qFAPI0j4GSsh7M0Q3Ypjzr4+CEVg82V8= +github.com/go-text/typesetting v0.2.1/go.mod h1:mTOxEwasOFpAMBjEQDhdWRckoLLeI/+qrQeBCTGEt6M= +github.com/go-text/typesetting-utils v0.0.0-20241103174707-87a29e9e6066 h1:qCuYC+94v2xrb1PoS4NIDe7DGYtLnU2wWiQe9a1B1c0= +github.com/go-text/typesetting-utils v0.0.0-20241103174707-87a29e9e6066/go.mod h1:DDxDdQEnB70R8owOx3LVpEFvpMK9eeH1o2r0yZhFI9o= +github.com/godbus/dbus/v5 v5.1.0 h1:4KLkAxT3aOY8Li4FRJe/KvhoNFFxo0m6fNuFUO8QJUk= +github.com/godbus/dbus/v5 v5.1.0/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3 h1:LMLX+LgTNWpfvCBdFebv6EsYotImrt/Ppc5cXIriCSo= github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3/go.mod h1:jl5iWTm0/hd5PjEYEOuwAJ57L/CibdZfrqZ5XA5GrCk= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/hack-pad/go-indexeddb v0.3.2 h1:DTqeJJYc1usa45Q5r52t01KhvlSN02+Oq+tQbSBI91A= +github.com/hack-pad/go-indexeddb v0.3.2/go.mod h1:QvfTevpDVlkfomY498LhstjwbPW6QC4VC/lxYb0Kom0= +github.com/hack-pad/safejs v0.1.0 h1:qPS6vjreAqh2amUqj4WNG1zIw7qlRQJ9K10eDKMCnE8= +github.com/hack-pad/safejs v0.1.0/go.mod h1:HdS+bKF1NrE72VoXZeWzxFOVQVUSqZJAG0xNCnb+Tio= github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k= github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= +github.com/jeandeaual/go-locale v0.0.0-20241217141322-fcc2cadd6f08 h1:wMeVzrPO3mfHIWLZtDcSaGAe2I4PW9B/P5nMkRSwCAc= +github.com/jeandeaual/go-locale v0.0.0-20241217141322-fcc2cadd6f08/go.mod h1:ZDXo8KHryOWSIqnsb/CiDq7hQUYryCgdVnxbj8tDG7o= +github.com/jsummers/gobmp v0.0.0-20230614200233-a9de23ed2e25 h1:YLvr1eE6cdCqjOe972w/cYF+FjW34v27+9Vo5106B4M= +github.com/jsummers/gobmp v0.0.0-20230614200233-a9de23ed2e25/go.mod h1:kLgvv7o6UM+0QSf0QjAse3wReFDsb9qbZJdfexWlrQw= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/mattn/go-isatty v0.0.24 h1:tGZZoVgT/KiqK1c8ocVLeDS8BSWMRd47J3Lbz7vsReI= github.com/mattn/go-isatty v0.0.24/go.mod h1:nMCL3Zebbrt45jsMDgnfIwz6ydEQApk5oEI3HqDio6A= github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w= github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls= +github.com/nfnt/resize v0.0.0-20180221191011-83c6a9932646 h1:zYyBkD/k9seD2A7fsi6Oo2LfFZAehjjQMERAvZLEDnQ= +github.com/nfnt/resize v0.0.0-20180221191011-83c6a9932646/go.mod h1:jpp1/29i3P1S/RLdc7JQKbRpFeM1dOBd8T9ki5s+AY8= +github.com/nicksnyder/go-i18n/v2 v2.5.1 h1:IxtPxYsR9Gp60cGXjfuR/llTqV8aYMsC472zD0D1vHk= +github.com/nicksnyder/go-i18n/v2 v2.5.1/go.mod h1:DrhgsSDZxoAfvVrBVLXoxZn/pN5TXqaDbq7ju94viiQ= +github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e h1:fD57ERR4JtEqsWbfPhv4DMiApHyliiK5xCTNVSPiaAs= +github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLANZcx1PVRCS0qkT7pwLkGfwJo4zjcN/Tysno= +github.com/pkg/profile v1.7.0 h1:hnbDkaNWPCLMO9wGLdBFTIZvzDrDfBM2072E1S9gJkA= +github.com/pkg/profile v1.7.0/go.mod h1:8Uer0jas47ZQMJ7VD+OHknK4YDY07LPUC6dEvqDjvNo= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE= github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo= +github.com/rymdport/portal v0.4.1 h1:2dnZhjf5uEaeDjeF/yBIeeRo6pNI2QAKm7kq1w/kbnA= +github.com/rymdport/portal v0.4.1/go.mod h1:kFF4jslnJ8pD5uCi17brj/ODlfIidOxlgUDTO5ncnC4= +github.com/srwiley/oksvg v0.0.0-20221011165216-be6e8873101c h1:km8GpoQut05eY3GiYWEedbTT0qnSxrCjsVbb7yKY1KE= +github.com/srwiley/oksvg v0.0.0-20221011165216-be6e8873101c/go.mod h1:cNQ3dwVJtS5Hmnjxy6AgTPd0Inb3pW05ftPSX7NZO7Q= +github.com/srwiley/rasterx v0.0.0-20220730225603-2ab79fcdd4ef h1:Ch6Q+AZUxDBCVqdkI8FSpFyZDtCVBc2VmejdNrm5rRQ= +github.com/srwiley/rasterx v0.0.0-20220730225603-2ab79fcdd4ef/go.mod h1:nXTWP6+gD5+LUJ8krVhhoeHjvHTutPxMYl5SvkcnJNE= +github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= +github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= +github.com/yuin/goldmark v1.7.8 h1:iERMLn0/QJeHFhxSt3p6PeN9mGnvIKSpG9YYorDMnic= +github.com/yuin/goldmark v1.7.8/go.mod h1:uzxRWxtg69N339t3louHJ7+O03ezfj6PlliRlaOzY1E= +golang.org/x/image v0.24.0 h1:AN7zRgVsbvmTfNyqIbbOraYL8mSwcKncEj8ofjgzcMQ= +golang.org/x/image v0.24.0/go.mod h1:4b/ITuLfqYq1hqZcjofwctIhi7sZh2WaCjvsBNjjya8= golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ= golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0= +golang.org/x/net v0.35.0 h1:T5GQRQb2y08kTAByq9L4/bz8cipCdA8FbRTXewonqY8= +golang.org/x/net v0.35.0/go.mod h1:EglIi67kWsHKlRzzVMUD93VMSWGFOMSZgxFjparz1Qk= golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM= golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/text v0.22.0 h1:bofq7m3/HAFvbF51jz3Q9wLg3jkvSPuiZu/pD1XwgtM= +golang.org/x/text v0.22.0/go.mod h1:YRoo4H8PVmsu+E3Ou7cqLVH8oXWIHVoX0jqUWALQhfY= golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q= golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA= -gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f h1:BLraFXnmrev5lT+xlilqcH8XK9/i0At2xKjWk4p6zsU= +gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= modernc.org/cc/v4 v4.29.1 h1:MKgdCV3WykTSPqpVrnxdEDS0HEd2FHpKZDzxzU5LyeI= diff --git a/internal/app/app.go b/internal/app/app.go index de7ffdb..e81bfbe 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -23,6 +23,7 @@ import ( "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/ui" "github.com/kamelion/ratatoskr-go/internal/update" "github.com/kamelion/ratatoskr-go/internal/worker" ) @@ -67,6 +68,11 @@ type App struct { // Events — доменная шина UI; LogEvents — шина логов (панель «Логи»). Events *events.Bus LogEvents *events.LogBus + + // UI — десктопное окно (если собрано с cgo и не задан --noui). + // nil — headless-режим (без окна). + UI ui.Window + Controller *ui.Controller } // New читает конфиг и собирает все зависимости. @@ -74,7 +80,7 @@ type App struct { // version — вшитая версия бинаря (ldflag -X main.version). // updateToken — вшитый токен read:package для авто-обновления // (ldflag -X main.updateToken); имеет приоритет над update.token из конфига. -func New(configPath, version, updateToken string) (*App, error) { +func New(configPath, version, updateToken string, noUI bool) (*App, error) { cfg, err := config.Load(configPath) if err != nil { return nil, fmt.Errorf("%w: %v", ErrConfig, err) @@ -165,6 +171,22 @@ func New(configPath, version, updateToken string) (*App, error) { router := chat.NewRouter(a.handleIncoming) + // UI-канал: ещё одна реализация chat.Channel (спец 12.1). + // Собирается только при наличии cgo и отсутствии --noui. + if !noUI { + win := newUIWindow(ui.NewDBStore(store)) + if win != nil { + if err := router.Attach(win); err != nil { + store.Close() + return nil, fmt.Errorf("attach ui: %w", err) + } + a.UI = win + // Контроллер подписывается на шины и мостит события в окно (fyne.Do). + a.Controller = ui.New(a.Events, a.LogEvents, win) + a.Controller.Start() + } + } + // Telegram-канал tg := telegram.New(cfg.Telegram.Token, cfg.Chat.PollInterval.Duration()) a.tg = tg @@ -246,6 +268,14 @@ func (a *App) Run(ctx context.Context) error { // Авто-проверка обновления (только уведомление владельца; замена — по /update) a.startAutoCheck(ctx) + // UI: окно блокирует главную горутину (спец 12.6). Кнопка «Завершить» + // отменяет ctx → fyneApp.Quit() → Run возвращается → graceful shutdown. + // Закрытие крестиком = сворачивание: Core продолжает работать. + if a.UI != nil { + a.UI.SetOnQuit(cancel) + return a.UI.Run(ctx) + } + // Ожидание сигнала или фатальной ошибки Telegram sigCh := make(chan os.Signal, 1) signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM) diff --git a/internal/app/app_test.go b/internal/app/app_test.go index e2f4d06..832caf7 100644 --- a/internal/app/app_test.go +++ b/internal/app/app_test.go @@ -21,7 +21,7 @@ func TestNew(t *testing.T) { t.Setenv("RATATOSKR_DB", dbPath) // Загружаем без config-файла (дефолты + env) - a, err := New("", "dev", "") + a, err := New("", "dev", "", true) if err != nil { t.Fatalf("New() err = %v", err) } @@ -52,7 +52,7 @@ func TestNew_MissingToken(t *testing.T) { t.Setenv("TG_CHAT_ID", "12345") t.Setenv("RATATOSKR_DB", dbPath) - _, err := New("", "dev", "") + _, err := New("", "dev", "", true) if err == nil { t.Fatal("expected error for missing token") } @@ -67,7 +67,7 @@ func TestNew_BadDB(t *testing.T) { t.Setenv("TG_CHAT_ID", "12345") t.Setenv("RATATOSKR_DB", dbPath) - _, err := New("", "dev", "") + _, err := New("", "dev", "", true) if err == nil { t.Fatal("expected error for invalid db path") } @@ -84,7 +84,7 @@ func TestNew_RunCtxCancel(t *testing.T) { t.Setenv("TG_CHAT_ID", "12345") t.Setenv("RATATOSKR_DB", dbPath) - a, err := New("", "dev", "") + a, err := New("", "dev", "", true) if err != nil { t.Fatalf("New() err = %v", err) } @@ -130,7 +130,7 @@ func TestNew_WorktreeCreated(t *testing.T) { t.Fatalf("write config: %v", err) } - a, err := New(configPath, "dev", "") + a, err := New(configPath, "dev", "", true) if err != nil { t.Fatalf("New() err = %v", err) } @@ -175,7 +175,7 @@ func TestNew_UpdateWiring(t *testing.T) { } // 1) base_url из update-блока, а НЕ git - a, err := New(configPath, "dev", "") + a, err := New(configPath, "dev", "", true) if err != nil { t.Fatalf("New() err = %v", err) } @@ -196,7 +196,7 @@ func TestNew_UpdateWiring(t *testing.T) { } // 2) вшитый updateToken перекрывает конфиг - a2, err := New(configPath, "dev", "embedded-update-token") + a2, err := New(configPath, "dev", "embedded-update-token", true) if err != nil { t.Fatalf("New() err = %v", err) } @@ -248,4 +248,4 @@ func TestVersion_Semver(t *testing.T) { if !re.MatchString(Version) { t.Errorf("Version = %q, ожидался формат major.minor.patch", Version) } -} \ No newline at end of file +} diff --git a/internal/app/ui_cgo.go b/internal/app/ui_cgo.go new file mode 100644 index 0000000..a7b8535 --- /dev/null +++ b/internal/app/ui_cgo.go @@ -0,0 +1,15 @@ +//go:build cgo + +package app + +import ( + "github.com/kamelion/ratatoskr-go/internal/ui" + "github.com/kamelion/ratatoskr-go/internal/ui/desktop" +) + +// newUIWindow создаёт десктопное окно (Fyne) на основе чтения-модели. +// Сборка с cgo: GLFW доступен, окно реально. +func newUIWindow(store ui.Store) ui.Window { + w := desktop.New("ratatoskr-go", store) + return w +} diff --git a/internal/app/ui_noui.go b/internal/app/ui_noui.go new file mode 100644 index 0000000..cd835e9 --- /dev/null +++ b/internal/app/ui_noui.go @@ -0,0 +1,11 @@ +//go:build !cgo + +package app + +import ( + "github.com/kamelion/ratatoskr-go/internal/ui" +) + +// newUIWindow возвращает nil: без cgo (нет GLFW) окно недоступно, +// приложение работает в headless-режиме. +func newUIWindow(ui.Store) ui.Window { return nil } diff --git a/internal/ui/commands.go b/internal/ui/commands.go new file mode 100644 index 0000000..bcf202e --- /dev/null +++ b/internal/ui/commands.go @@ -0,0 +1,46 @@ +package ui + +import "fmt" + +// Commands — отображение действий UI в текстовые команды канала. +// +// UI — это chat.Channel, поэтому все действия пользователя (кнопки, ввод) +// превращаются в текстовые сообщения, которые уходят в Router → Core. Это +// сохраняет единый путь команд (спец 12.2): никаких прямых вызовов Core из +// виджетов — только текст через канал. +type Commands struct { + // Submit отправляет текст пользователя в канал (обычно — window.Submit). + Submit func(text string) +} + +// NewCommands создаёт Commands с отправкой через fn. +func NewCommands(fn func(text string)) *Commands { + if fn == nil { + fn = func(string) {} + } + return &Commands{Submit: fn} +} + +// Start — новая задача (/start). +func (c *Commands) Start() { c.Submit("/start") } + +// Cancel — отмена текущей задачи (/cancel). +func (c *Commands) Cancel() { c.Submit("/cancel") } + +// Skip — пропустить сбор, сформировать черновик (/skip). +func (c *Commands) Skip() { c.Submit("/skip") } + +// Retry — перезапустить задачу N (/retry N). +func (c *Commands) Retry(id int64) { c.Submit(fmt.Sprintf("/retry %d", id)) } + +// Status — запросить статус задачи N (/status N). +func (c *Commands) Status(id int64) { c.Submit(fmt.Sprintf("/status %d", id)) } + +// Continue — продолжить задачу N (/continue N). +func (c *Commands) Continue(id int64) { c.Submit(fmt.Sprintf("/continue %d", id)) } + +// Approve — одобрить черновик («создавай»). +func (c *Commands) Approve() { c.Submit("создавай") } + +// SendText — обычное сообщение пользователя (ввод в поле). +func (c *Commands) SendText(text string) { c.Submit(text) } diff --git a/internal/ui/controller.go b/internal/ui/controller.go index b53f606..3765746 100644 --- a/internal/ui/controller.go +++ b/internal/ui/controller.go @@ -49,4 +49,4 @@ func (c *Controller) Start() { func (c *Controller) Close() { c.domain.Close() c.logs.Close() -} \ No newline at end of file +} diff --git a/internal/ui/controller_test.go b/internal/ui/controller_test.go index 67df3f9..9280f81 100644 --- a/internal/ui/controller_test.go +++ b/internal/ui/controller_test.go @@ -143,4 +143,4 @@ func TestControllerIdempotentClose(t *testing.T) { c.Start() c.Close() c.Close() // не должно падать -} \ No newline at end of file +} diff --git a/internal/ui/desktop/window.go b/internal/ui/desktop/window.go new file mode 100644 index 0000000..cd21481 --- /dev/null +++ b/internal/ui/desktop/window.go @@ -0,0 +1,497 @@ +//go:build cgo + +// Package desktop — Fyne-реализация окна UI (см. internal/ui). +// +// Сборка требует cgo (GLFW). На машинах без C-компилятора пакет не входит +// в сборку, и app работает в headless-режиме (newUIWindow → nil). +package desktop + +import ( + "context" + "log" + "strings" + "sync" + + "fyne.io/fyne/v2" + "fyne.io/fyne/v2/app" + "fyne.io/fyne/v2/container" + "fyne.io/fyne/v2/theme" + "fyne.io/fyne/v2/widget" + + "github.com/kamelion/ratatoskr-go/internal/chat" + "github.com/kamelion/ratatoskr-go/internal/events" + "github.com/kamelion/ratatoskr-go/internal/model" + "github.com/kamelion/ratatoskr-go/internal/storage" + "github.com/kamelion/ratatoskr-go/internal/ui" +) + +const ( + maxLogLen = 200_000 // обрезка буферов панелей, чтобы не расти бесконечно + selectedPref = "task.selected" + splitHPref = "split.h" + splitVPref = "split.v" +) + +// Window — Fyne-окно: chat.Channel + ui.View. +type Window struct { + fyneApp fyne.App + win fyne.Window + store ui.Store + commands *ui.Commands + + handler chat.Handler + + // состояние списка + mu sync.Mutex // защищает tasks + tasks []storage.Task + selected int64 + + // виджеты + list *widget.List + titleLbl *widget.Label + statusLbl *widget.Label + detailLbl *widget.Label + convLbl *widget.Label + logsLbl *widget.Label + stateLbl *widget.Label + input *widget.Entry + + hsplit *container.Split + vsplit *container.Split + + onQuit func() +} + +// New создаёт окно. appID — идентификатор Fyne-приложения (Preferences). +func New(appID string, store ui.Store) *Window { + fa := app.NewWithID(appID) + w := &Window{ + fyneApp: fa, + store: store, + } + w.win = fa.NewWindow("Ratatoskr") + w.win.SetFullScreen(true) + w.commands = ui.NewCommands(w.submitText) + w.build() + w.win.SetCloseIntercept(func() { + // Закрытие окна = сворачивание (спец 12.6): Core продолжает работать. + w.saveLayout() + w.win.Hide() + }) + fa.Settings().SetTheme(theme.DarkTheme()) + return w +} + +// SetOnQuit задаёт колбэк полного выхода (кнопка «Завершить»). +func (w *Window) SetOnQuit(fn func()) { w.onQuit = fn } + +// build собирает виджеты и раскладку. +func (w *Window) build() { + w.titleLbl = widget.NewLabel("—") + w.titleLbl.TextStyle = fyne.TextStyle{Bold: true} + w.statusLbl = widget.NewLabel("") + w.statusLbl.TextStyle = fyne.TextStyle{Italic: true} + w.detailLbl = widget.NewLabel("") + w.detailLbl.Wrapping = fyne.TextWrapWord + + w.convLbl = widget.NewLabel("Выберите задачу.") + w.convLbl.Wrapping = fyne.TextWrapWord + + w.logsLbl = widget.NewLabel("") + w.logsLbl.Wrapping = fyne.TextWrapWord + + w.stateLbl = widget.NewLabel("") + w.stateLbl.Wrapping = fyne.TextWrapWord + + // Список задач (слева) + w.list = widget.NewList( + func() int { w.mu.Lock(); defer w.mu.Unlock(); return len(w.tasks) }, + func() fyne.CanvasObject { + return widget.NewLabel("loading…") + }, + func(id widget.ListItemID, o fyne.CanvasObject) { + w.mu.Lock() + defer w.mu.Unlock() + if id < 0 || id >= len(w.tasks) { + return + } + t := w.tasks[id] + o.(*widget.Label).SetText(taskTitle(t)) + }, + ) + w.list.OnSelected = func(id widget.ListItemID) { + w.mu.Lock() + if id < 0 || id >= len(w.tasks) { + w.mu.Unlock() + return + } + t := w.tasks[id] + w.mu.Unlock() + w.selectTask(t.ID) + } + + quitBtn := widget.NewButtonWithIcon("Завершить", theme.LogoutIcon(), func() { + w.saveLayout() + if w.onQuit != nil { + w.onQuit() + } + }) + left := container.NewBorder(nil, quitBtn, nil, nil, w.list) + + // Рабочая область: детали сверху, диалог снизу (сплит 2×2 в плане). + details := container.NewVBox(w.titleLbl, w.statusLbl, w.detailLbl) + convScroll := container.NewScroll(w.convLbl) + right := container.NewVSplit( + container.NewScroll(details), + convScroll, + ) + w.vsplit = right + + // Нижняя панель: вкладки «Логи» + «Состояние». + bottomTabs := container.NewAppTabs( + container.NewTabItem("Логи", container.NewScroll(w.logsLbl)), + container.NewTabItem("Состояние", container.NewScroll(w.stateLbl)), + ) + + // Ввод + кнопки команд + w.input = widget.NewEntry() + w.input.SetPlaceHolder("Сообщение… (Enter — отправить)") + w.input.OnSubmitted = func(s string) { + s = strings.TrimSpace(s) + if s == "" { + return + } + w.input.SetText("") + w.commands.SendText(s) + } + cmdBar := container.NewHBox( + widget.NewButton("Новая", func() { w.commands.Start() }), + widget.NewButton("Создавай", func() { w.commands.Approve() }), + widget.NewButton("Пропустить", func() { w.commands.Skip() }), + widget.NewButton("Отмена", func() { w.commands.Cancel() }), + ) + inputRow := container.NewBorder(nil, nil, cmdBar, w.input, nil) + bottom := container.NewBorder(nil, inputRow, nil, nil, bottomTabs) + + w.hsplit = container.NewHSplit(left, right) + + root := container.NewVSplit(w.hsplit, bottom) + w.win.SetContent(root) + + // Восстановление layout из Preferences. + pref := w.fyneApp.Preferences() + w.hsplit.SetOffset(pref.FloatWithFallback(splitHPref, 0.30)) + w.vsplit.SetOffset(pref.FloatWithFallback(splitVPref, 0.45)) + + // Начальный снимок списка. + w.refreshList() +} + +// saveLayout сохраняет текущие позиции сплитов в Preferences. +// Вызывается при сворачивании окна и перед полным выходом. +func (w *Window) saveLayout() { + pref := w.fyneApp.Preferences() + pref.SetFloat(splitHPref, w.hsplit.Offset) + pref.SetFloat(splitVPref, w.vsplit.Offset) +} + +// taskTitle — строка в списке задач. +func taskTitle(t storage.Task) string { + title := t.Title + if title == "" { + title = "(без названия)" + } + return title + "\n" + string(t.Status) +} + +// selectTask загружает детали выбранной задачи (снимки из Store). +func (w *Window) selectTask(id int64) { + w.selected = id + w.fyneApp.Preferences().SetInt(selectedPref, int(id)) + + ctx := context.Background() + t, err := w.store.GetTask(ctx, id) + if err != nil { + w.titleLbl.SetText("—") + w.statusLbl.SetText("") + w.detailLbl.SetText("") + w.convLbl.SetText("") + return + } + w.titleLbl.SetText(t.Title) + w.statusLbl.SetText(string(t.Status)) + detail := "Цель: " + t.Goal + if len(t.Repos) > 0 { + detail += "\nРепозитории: " + strings.Join(t.Repos, ", ") + } + w.detailLbl.SetText(detail) + + // Диалог: история + краткие трассы. + hist, err := w.store.GetHistory(ctx, id) + if err != nil { + hist = nil + } + var b strings.Builder + for _, h := range hist { + b.WriteString(formatRole(h.Role)) + b.WriteString(h.Content) + b.WriteString("\n\n") + } + traces, err := w.store.GetTraces(ctx, id) + if err == nil { + for _, tr := range traces { + b.WriteString("— " + tr.Agent + " (" + string(tr.Status) + ")\n") + } + } + if b.Len() == 0 { + b.WriteString("Нет сообщений. /start — начать задачу.") + } + w.convLbl.SetText(b.String()) + + // Состояние: сброс к снимку задач. + w.refreshState(ctx, id) +} + +func formatRole(role string) string { + switch role { + case "user": + return "👤 " + default: + return "🤖 " + } +} + +// refreshState — снимок «Состояния» выбранной задачи (агенты + трейсы). +func (w *Window) refreshState(ctx context.Context, id int64) { + traces, err := w.store.GetTraces(ctx, id) + if err != nil { + return + } + var b strings.Builder + for _, tr := range traces { + b.WriteString(strings.ToUpper(tr.Agent) + ": " + string(tr.Status)) + if tr.SessionID != "" { + b.WriteString(" (session " + tr.SessionID + ")") + } + b.WriteString("\n") + } + if b.Len() == 0 { + b.WriteString("Нет активных агентов.") + } + w.stateLbl.SetText(b.String()) +} + +// refreshList перечитывает список задач из Store (снимок). +func (w *Window) refreshList() { + ctx := context.Background() + tasks, err := w.store.ListTasks(ctx) + if err != nil { + log.Printf("ui: list tasks: %v", err) + return + } + w.mu.Lock() + w.tasks = tasks + selected := w.selected + w.mu.Unlock() + + if w.list != nil { + w.list.Refresh() + } + _ = selected +} + +// appendConv добавляет строку в диалог (буфер ограничен). +func (w *Window) appendConv(text string) { + s := w.convLbl.Text + text + "\n\n" + if len(s) > maxLogLen { + s = s[len(s)-maxLogLen:] + } + w.convLbl.SetText(s) +} + +// appendLog добавляет строку в панель «Логи». +func (w *Window) appendLog(text string) { + s := w.logsLbl.Text + text + "\n" + if len(s) > maxLogLen { + s = s[len(s)-maxLogLen:] + } + w.logsLbl.SetText(s) +} + +// appendState добавляет строку в панель «Состояние». +func (w *Window) appendState(text string) { + s := w.stateLbl.Text + text + "\n" + if len(s) > maxLogLen { + s = s[len(s)-maxLogLen:] + } + w.stateLbl.SetText(s) +} + +// submitText отправляет ввод пользователя в Router. +func (w *Window) submitText(text string) { + if w.handler != nil { + w.handler(chat.Incoming{ + UserID: ui.UID, + Address: ui.Address, + Channel: w, + Msg: chat.Message{Text: text}, + }) + } +} + +// ---- chat.Channel ---- + +// OnMessage регистрирует обработчик входящих (Router). +func (w *Window) OnMessage(h chat.Handler) { w.handler = h } + +// Run показывает окно и запускает Fyne-цикл. Блокирует до Quit. +func (w *Window) Run(ctx context.Context) error { + w.win.Show() + go func() { + <-ctx.Done() + w.fyneApp.Quit() + }() + w.fyneApp.Run() + return nil +} + +// Send доставляет сообщение Core в диалог окна. +func (w *Window) Send(_ context.Context, _ chat.Address, m chat.Message) error { + text := m.Text + if len(m.Options) > 0 { + var sb strings.Builder + sb.WriteString(text) + for i, opt := range m.Options { + sb.WriteString("\n") + sb.WriteString(opt.Label) + _ = i + } + text = sb.String() + } + fyne.Do(func() { w.appendConv("🤖 " + text) }) + return nil +} + +// Ask — вопрос с вариантами (отображается в диалоге). +func (w *Window) Ask(_ context.Context, _ chat.Address, m chat.Message) error { + text := m.Text + for _, opt := range m.Options { + text += "\n" + opt.Label + } + fyne.Do(func() { w.appendConv("❓ " + text) }) + return nil +} + +// Close завершает Fyne-цикл. +func (w *Window) Close() error { + w.fyneApp.Quit() + return nil +} + +// ---- ui.View ---- + +// Все колбэки выполняются в горутинах events.Hub → обновляем UI через fyne.Do. + +func (w *Window) OnTaskCreated(e events.TaskCreated) { + fyne.Do(func() { + w.refreshList() + if w.selected == 0 { + w.selectTask(e.ID) + } + }) +} + +func (w *Window) OnTaskUpdated(e events.TaskUpdated) { + fyne.Do(func() { + w.refreshList() + if w.selected == e.ID { + w.selectTask(e.ID) + } + }) +} + +func (w *Window) OnTaskDeleted(e events.TaskDeleted) { + fyne.Do(func() { + if w.selected == e.ID { + w.selected = 0 + } + w.refreshList() + }) +} + +func (w *Window) OnTaskStatusChanged(e events.TaskStatusChanged) { + fyne.Do(func() { + w.refreshList() + if w.selected == e.ID { + w.selectTask(e.ID) + } + w.appendState(taskStatusLine(e)) + }) +} + +func (w *Window) OnHistoryAppended(e events.HistoryAppended) { + fyne.Do(func() { + if w.selected == e.TaskID { + w.appendConv(formatRole(e.Role) + e.Content) + } + }) +} + +func (w *Window) OnTraceAppended(e events.TraceAppended) { + fyne.Do(func() { + w.appendState(strings.ToUpper(e.Agent) + ": " + string(e.Status)) + }) +} + +func (w *Window) OnAgentActivity(e events.AgentActivity) { + fyne.Do(func() { + w.appendState(strings.ToUpper(e.Agent) + " → " + e.Stage) + }) +} + +func (w *Window) OnLog(e events.LogLine) { + fyne.Do(func() { w.appendLog(e.Text) }) +} + +func taskStatusLine(e events.TaskStatusChanged) string { + return statusBadge(e.To) + " #" + itoa(e.ID) + ": " + string(e.From) + " → " + string(e.To) +} + +func statusBadge(s model.Status) string { + switch s { + case model.StatusSuccess: + return "✅" + case model.StatusFailed, model.StatusAborted: + return "❌" + case model.StatusRunning, model.StatusCollecting: + return "⏳" + case model.StatusReady, model.StatusApproved: + return "🟡" + case model.StatusCancelled: + return "🚫" + default: + return "•" + } +} + +func itoa(v int64) string { + if v == 0 { + return "0" + } + neg := v < 0 + if neg { + v = -v + } + var b [24]byte + i := len(b) + for v > 0 { + i-- + b[i] = byte('0' + v%10) + v /= 10 + } + if neg { + i-- + b[i] = '-' + } + return string(b[i:]) +} diff --git a/internal/ui/store.go b/internal/ui/store.go new file mode 100644 index 0000000..183c310 --- /dev/null +++ b/internal/ui/store.go @@ -0,0 +1,73 @@ +package ui + +import ( + "context" + + "github.com/kamelion/ratatoskr-go/internal/storage" +) + +// Store — абстракция чтения-модели для UI (спец 12.5). +// +// Возвращает только копии (value-типы), никогда — живые указатели на +// разделяемые контейнеры core/БД. UI читает снимки и реагирует на события; +// мутирует состояние только Core. +type Store interface { + // ListTasks возвращает все задачи (копии), сортировка по updated_at desc. + ListTasks(ctx context.Context) ([]storage.Task, error) + // GetTask возвращает задачу по ID (копию). + GetTask(ctx context.Context, id int64) (storage.Task, error) + // GetHistory возвращает историю диалога задачи (копии). + GetHistory(ctx context.Context, id int64) ([]storage.HistoryMsg, error) + // GetTraces возвращает трассы задачи (копии). + GetTraces(ctx context.Context, id int64) ([]storage.Trace, error) +} + +// DBStore — реализация Store поверх *storage.Storage (БД). +type DBStore struct { + s *storage.Storage +} + +// NewDBStore создаёт DBStore на основе БД. +func NewDBStore(s *storage.Storage) *DBStore { + return &DBStore{s: s} +} + +// ListTasks реализует Store: deref-копии записей, без указателей. +func (d *DBStore) ListTasks(ctx context.Context) ([]storage.Task, error) { + rows, err := d.s.ListTasks(ctx, storage.TaskFilter{Limit: 200}) + if err != nil { + return nil, err + } + out := make([]storage.Task, 0, len(rows)) + for _, r := range rows { + out = append(out, *r) + } + return out, nil +} + +// GetTask реализует Store. +func (d *DBStore) GetTask(ctx context.Context, id int64) (storage.Task, error) { + t, err := d.s.GetTask(ctx, id) + if err != nil { + return storage.Task{}, err + } + return *t, nil +} + +// GetHistory реализует Store. +func (d *DBStore) GetHistory(ctx context.Context, id int64) ([]storage.HistoryMsg, error) { + return d.s.GetHistory(ctx, id) +} + +// GetTraces реализует Store. +func (d *DBStore) GetTraces(ctx context.Context, id int64) ([]storage.Trace, error) { + rows, err := d.s.GetTraces(ctx, id) + if err != nil { + return nil, err + } + out := make([]storage.Trace, 0, len(rows)) + for _, r := range rows { + out = append(out, *r) + } + return out, nil +} diff --git a/internal/ui/store_test.go b/internal/ui/store_test.go new file mode 100644 index 0000000..e668e9c --- /dev/null +++ b/internal/ui/store_test.go @@ -0,0 +1,107 @@ +package ui + +import ( + "context" + "testing" + + "github.com/kamelion/ratatoskr-go/internal/storage" +) + +func newTestStore(t *testing.T) *storage.Storage { + t.Helper() + s, err := storage.Open(context.Background(), ":memory:") + if err != nil { + t.Fatalf("open storage: %v", err) + } + t.Cleanup(func() { s.Close() }) + return s +} + +func TestDBStoreListTasksReturnsCopies(t *testing.T) { + s := newTestStore(t) + ctx := context.Background() + + id1, err := s.CreateTask(ctx, &storage.Task{ChatID: "ui://local", Title: "A"}) + if err != nil { + t.Fatalf("create: %v", err) + } + if _, err := s.CreateTask(ctx, &storage.Task{ChatID: "tg://1", Title: "B"}); err != nil { + t.Fatalf("create: %v", err) + } + + st := NewDBStore(s) + tasks, err := st.ListTasks(ctx) + if err != nil { + t.Fatalf("ListTasks: %v", err) + } + if len(tasks) != 2 { + t.Fatalf("got %d tasks, want 2", len(tasks)) + } + + // мутация полученной копии не влияет на БД + tasks[0].Title = "mutated" + fetched, err := s.GetTask(ctx, id1) + if err != nil { + t.Fatalf("get: %v", err) + } + if fetched.Title == "mutated" { + t.Fatal("mutating snapshot leaked into DB") + } +} + +func TestDBStoreGetTaskAndHistory(t *testing.T) { + s := newTestStore(t) + ctx := context.Background() + + id, err := s.CreateTask(ctx, &storage.Task{ChatID: "ui://local", Title: "A"}) + if err != nil { + t.Fatalf("create: %v", err) + } + if err := s.AppendHistory(ctx, id, "user", "hello"); err != nil { + t.Fatalf("append: %v", err) + } + + st := NewDBStore(s) + task, err := st.GetTask(ctx, id) + if err != nil { + t.Fatalf("GetTask: %v", err) + } + if task.ID != id { + t.Fatalf("task.ID = %d, want %d", task.ID, id) + } + + hist, err := st.GetHistory(ctx, id) + if err != nil { + t.Fatalf("GetHistory: %v", err) + } + if len(hist) != 1 || hist[0].Content != "hello" { + t.Fatalf("bad history: %+v", hist) + } +} + +func TestDBStoreGetTraces(t *testing.T) { + s := newTestStore(t) + ctx := context.Background() + + id, err := s.CreateTask(ctx, &storage.Task{ChatID: "ui://local"}) + if err != nil { + t.Fatalf("create: %v", err) + } + tid, err := s.AppendTrace(ctx, &storage.Trace{TaskID: id, Agent: "dev", Prompt: "p"}) + if err != nil { + t.Fatalf("append trace: %v", err) + } + _ = s.UpdateTraceStatus(ctx, tid, storage.TraceSuccess) + + st := NewDBStore(s) + traces, err := st.GetTraces(ctx, id) + if err != nil { + t.Fatalf("GetTraces: %v", err) + } + if len(traces) != 1 { + t.Fatalf("got %d traces, want 1", len(traces)) + } + if traces[0].Status != storage.TraceSuccess { + t.Fatalf("trace status = %s, want success", traces[0].Status) + } +} diff --git a/internal/ui/view.go b/internal/ui/view.go index 38df43c..da4e261 100644 --- a/internal/ui/view.go +++ b/internal/ui/view.go @@ -37,4 +37,4 @@ func (NilView) OnTaskStatusChanged(events.TaskStatusChanged) {} func (NilView) OnHistoryAppended(events.HistoryAppended) {} func (NilView) OnTraceAppended(events.TraceAppended) {} func (NilView) OnAgentActivity(events.AgentActivity) {} -func (NilView) OnLog(events.LogLine) {} \ No newline at end of file +func (NilView) OnLog(events.LogLine) {} diff --git a/internal/ui/window.go b/internal/ui/window.go new file mode 100644 index 0000000..9850fc2 --- /dev/null +++ b/internal/ui/window.go @@ -0,0 +1,68 @@ +package ui + +import ( + "context" + + "github.com/kamelion/ratatoskr-go/internal/chat" + "github.com/kamelion/ratatoskr-go/internal/events" +) + +// UID — фиксированный пользовательский идентификатор UI-канала (один пользователь). +const UID = chat.UserID("ui://local") + +// Address — адрес доставки для UI-канала. +const Address = chat.Address("ui://local") + +// Window — окно десктопного UI. +// +// Объединяет две роли: +// - chat.Channel — ещё одна реализация канала (спец 12.1): входящие от +// пользователя (ввод) уходят в Router → Core; исходящие (Send/Ask) — из +// Core/Router попадают в окно (диалог). +// - View — событийный мост из Controller (статусы, история, трейсы, логи). +// +// Реализация (Fyne) живёт в internal/ui/desktop под build-tag `cgo`. +type Window interface { + chat.Channel + View + // SetOnQuit задаёт колбэк полного выхода (кнопка «Завершить»). + SetOnQuit(func()) +} + +// NilWindow — no-op окно для headless-режима (--noui) и тестов. +// Реализует Window: канал без доставки и View без действий. +type NilWindow struct { + handler chat.Handler +} + +// NewNilWindow создаёт NilWindow. +func NewNilWindow() *NilWindow { return &NilWindow{} } + +func (n *NilWindow) OnMessage(h chat.Handler) { n.handler = h } + +func (n *NilWindow) SetOnQuit(func()) {} + +func (n *NilWindow) Run(ctx context.Context) error { + <-ctx.Done() + return ctx.Err() +} + +func (n *NilWindow) Send(context.Context, chat.Address, chat.Message) error { return nil } +func (n *NilWindow) Ask(context.Context, chat.Address, chat.Message) error { return nil } +func (n *NilWindow) Close() error { return nil } + +// Submit отправляет текст пользователя в Router (ввод в окне). +func (n *NilWindow) Submit(text string) { + if n.handler != nil { + n.handler(chat.Incoming{UserID: UID, Address: Address, Channel: n, Msg: chat.Message{Text: text}}) + } +} + +func (n *NilWindow) OnTaskCreated(events.TaskCreated) {} +func (n *NilWindow) OnTaskUpdated(events.TaskUpdated) {} +func (n *NilWindow) OnTaskDeleted(events.TaskDeleted) {} +func (n *NilWindow) OnTaskStatusChanged(events.TaskStatusChanged) {} +func (n *NilWindow) OnHistoryAppended(events.HistoryAppended) {} +func (n *NilWindow) OnTraceAppended(events.TraceAppended) {} +func (n *NilWindow) OnAgentActivity(events.AgentActivity) {} +func (n *NilWindow) OnLog(events.LogLine) {} diff --git a/internal/ui/window_test.go b/internal/ui/window_test.go new file mode 100644 index 0000000..75ba9d4 --- /dev/null +++ b/internal/ui/window_test.go @@ -0,0 +1,80 @@ +package ui + +import ( + "context" + "testing" + "time" + + "github.com/kamelion/ratatoskr-go/internal/chat" +) + +func TestCommandsMapToText(t *testing.T) { + var got []string + c := NewCommands(func(text string) { got = append(got, text) }) + + c.Start() + c.Cancel() + c.Skip() + c.Retry(7) + c.Status(7) + c.Continue(3) + c.Approve() + c.SendText("просто текст") + + want := []string{"/start", "/cancel", "/skip", "/retry 7", "/status 7", "/continue 3", "создавай", "просто текст"} + if len(got) != len(want) { + t.Fatalf("got %d commands, want %d: %v", len(got), len(want), got) + } + for i := range want { + if got[i] != want[i] { + t.Errorf("command %d = %q, want %q", i, got[i], want[i]) + } + } +} + +func TestNilWindowSubmitReachesHandler(t *testing.T) { + w := NewNilWindow() + + incoming := make(chan chat.Incoming, 1) + w.OnMessage(func(inc chat.Incoming) { + incoming <- inc + }) + + w.Submit("/start") + + select { + case inc := <-incoming: + if inc.UserID != UID { + t.Errorf("UserID = %q, want %q", inc.UserID, UID) + } + if inc.Address != Address { + t.Errorf("Address = %q, want %q", inc.Address, Address) + } + if inc.Msg.Text != "/start" { + t.Errorf("Msg.Text = %q, want /start", inc.Msg.Text) + } + case <-time.After(2 * time.Second): + t.Fatal("handler not called") + } +} + +func TestNilWindowRunStopsOnCancel(t *testing.T) { + w := NewNilWindow() + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { done <- w.Run(ctx) }() + time.Sleep(50 * time.Millisecond) + cancel() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("Run did not stop on cancel") + } +} + +func TestNilWindowSendNoop(t *testing.T) { + w := NewNilWindow() + if err := w.Send(context.Background(), Address, chat.Message{Text: "hi"}); err != nil { + t.Fatalf("Send: %v", err) + } +} -- 2.49.1