15 Commits

Author SHA1 Message Date
ki.sagidullin
5c82229a00 chore(serena): обновить память проекта (onboarding)
Some checks failed
CI / test (pull_request) Failing after 1m13s
CI / build-and-package (amd64, linux) (pull_request) Failing after 1m2s
CI / build-and-package (amd64, windows) (pull_request) Successful in 28s
2026-08-21 00:08:45 +05:00
ki.sagidullin
612446d65b chore(scripts): add build-publish-ui.ps1 2026-08-21 00:05:54 +05:00
60a6966a95 Merge pull request 'feat/new_ui' (#8) from feat/new_ui into main
Some checks failed
CI / test (push) Failing after 1m8s
CI / build-and-package (amd64, linux) (push) Failing after 58s
CI / build-and-package (amd64, windows) (push) Successful in 28s
Reviewed-on: http://gitea.hal9000.home/kamelion/ratatoskr-go/pulls/8
2026-08-20 14:50:07 +05:00
ki.sagidullin
c2272137b3 feat(ui): команды, snapshots, Fyne-окно и интеграция (--noui)
Some checks failed
CI / test (pull_request) Failing after 1m19s
CI / build-and-package (amd64, linux) (pull_request) Failing after 58s
CI / build-and-package (amd64, windows) (pull_request) Successful in 53s
Команды (спец 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).
2026-08-20 08:54:55 +05:00
ki.sagidullin
baf2ad7147 feat(ui): контроллер — подписчик шин и мост к Fyne-представлению
- View — интерфейс Fyne-слоя: колбэки на все доменные события + логи.
- Controller подписывается на доменную и логовую шины (через events.Hub)
  и диспатчит события в View; колбэки выполняются в горутинах Hub, поэтому
  реализация Fyne обязана обновлять виджеты через fyne.Do (спец 12.4).
- Close отписывает контроллер (сворачивание окна) без остановки Core.
- NilView — no-op для headless (--noui) и тестов.
- Юнит-тесты: доставка доменных/логовых событий, остановка после Close.
2026-08-20 08:15:42 +05:00
ki.sagidullin
66e40d384e feat(events): типизированный Hub — подписчик шины с диспетчеризацией
- Hub слушает *Bus в собственной горутине и вызывает зарегистрированные
  обработчики по типу события (On[T]); порядок сохраняется.
- Мост к UI: колбэки выполняются в горутине Hub → внутри можно переложить
  работу на поток Fyne (fyne.Do) или thread-safe binding.
- Close останавливает горутину и отписывается; Start идемпотентен.
- Юнит-тесты: доставка по типу, игнор посторонних типов, несколько
  обработчиков одного типа, остановка после Close.
2026-08-20 07:17:57 +05:00
ki.sagidullin
4957a5552e feat(events): внедрить Publisher в Core/Worker/Analyst + шины в app
- Core: setStatus публикует TaskStatusChanged на каждом переходе
  (handleStart/Cancel/Retry/Consent/Edit/Turn, runDecide).
- Worker: setStatus публикует TaskStatusChanged (running/timeout/failed/
  success, failTask), notify остаётся авто-уведомлением.
- Analyst: Decide публикует AgentActivity (stage=decide).
- app: создаёт доменную (256) и логовую (1024) шины, подключает их к
  Core/Worker/Analyst, редиректит log в обе (stderr + логовая шина).
2026-08-20 04:30:29 +05:00
ki.sagidullin
1f7ab9a67e feat(events): EventBus + domain-модель статусов
- internal/model: единый источник статусов задачи (Status, TraceStatus,
  машина переходов), без зависимости от storage.
- internal/storage: совместимый мост (type Status = model.Status,
  re-export констант) — внешний код не меняется.
- internal/events: шина событий (fan-out, блокирующий Publish с гарантией
  порядка), события задач/трейсов, отдельная логовая шина + LogWriter,
  Publisher/NilPublisher для внедрения в Core/Worker.
- docs/ui-spec.md: спецификация десктопного UI (Fyne).
2026-08-20 01:02:31 +05:00
408d137747 Merge pull request 'fix(opencode): считать реальный прогресс стрима для idle-детекции' (#7) from feat/3c8528d into main
All checks were successful
CI / test (push) Successful in 1m3s
CI / build-and-package (amd64, linux) (push) Successful in 54s
CI / build-and-package (amd64, windows) (push) Successful in 53s
Reviewed-on: http://gitea.hal9000.home/kamelion/ratatoskr-go/pulls/7
2026-08-19 19:12:56 +05:00
ki.sagidullin
be4d749c45 fix(opencode): считать реальный прогресс стрима для idle-детекции
All checks were successful
CI / test (pull_request) Successful in 1m9s
CI / build-and-package (amd64, linux) (pull_request) Successful in 53s
CI / build-and-package (amd64, windows) (pull_request) Successful in 52s
Растущий в один text-парт стрим (text-delta) и reasoning больше не
выглядят как зависшая нейронка: idle-таймер сбрасывается по росту
числа партов и суммарной длины text/reasoning.
2026-08-19 19:11:26 +05:00
0b31e18d85 Merge pull request 'test: починить тесты на Windows' (#6) from feat/3c8528d into main
All checks were successful
CI / test (push) Successful in 35s
CI / build-and-package (amd64, linux) (push) Successful in 45s
CI / build-and-package (amd64, windows) (push) Successful in 42s
Reviewed-on: http://gitea.hal9000.home/kamelion/ratatoskr-go/pulls/6
2026-08-19 10:50:23 +05:00
ki.sagidullin
b5fb583c90 test: починить тесты на Windows
All checks were successful
CI / test (pull_request) Successful in 41s
CI / build-and-package (amd64, linux) (pull_request) Successful in 37s
CI / build-and-package (amd64, windows) (pull_request) Successful in 34s
- E2E (app): e2eFakeAPI переведён на v2 HTTP API opencode (/api/*) с
  определением агента по тексту промпта; Router получает processed-счётчик
  и WaitProcessed, e2eChannel.deliver ждёт асинхронную обработку — убирает
  гонку «запрос сразу после deliver» и коллатеральный 'database is closed'.
- app_test: одинарные YAML-кавычки для путей Windows (backslash-escape) +
  закрытие Store в TestNew/TestNew_RunCtxCancel/TestNew_UpdateWiring.
- config_test: абсолютный путь строится с корнем тома (C:\...) и одинарными
  кавычками YAML.
- opencode/server_test: fakeServeBin на Windows — .cmd с ping (#!/bin/sh
  не исполняется).
- worker_test: TestWorkerSemaphore поллит до целевого статуса вместо
  фиксированных sleep (git на Windows медленнее).
2026-08-19 10:47:55 +05:00
3c8528dbd9 Merge pull request 'feat(opencode): переход на v2 HTTP API opencode (хардпин модели, поллинг вердикта)' (#5) from feat/ad025c1 into main
Some checks failed
CI / test (push) Failing after 30s
CI / build-and-package (amd64, linux) (push) Successful in 49s
CI / build-and-package (amd64, windows) (push) Successful in 52s
Reviewed-on: http://gitea.hal9000.home/kamelion/ratatoskr-go/pulls/5
2026-08-19 08:55:16 +05:00
ki.sagidullin
9f64be4ea5 feat(opencode): переход на v2 HTTP API opencode (хардпин модели, поллинг вердикта)
Some checks failed
CI / test (pull_request) Failing after 32s
CI / build-and-package (amd64, linux) (pull_request) Successful in 42s
CI / build-and-package (amd64, windows) (pull_request) Successful in 42s
- client.go: эндпоинты /api/* (create+model, prompt-admit, message, active, interrupt)
- runner.go: неблокирующий prompt + поллинг новых assistant-сообщений;
  завершение = сессия ушла из активных дренажей + стабильное финальное сообщение
- config.go: чтение top-level model из opencode.jsonc (JSONC-стрип) + хардпин в сессию
- server.go: healthcheck /api/health, MinVersion=1.18.18, понятная ошибка для старого бинаря
- класс O5 WARN: устойчивость к v1-конфигу провайдера (npm/options игнорируются v2)
- README: раздел интеграции, минимальная версия opencode, предупреждения
- .serena: актуализация памяти (core, tech_stack)
2026-08-19 08:29:11 +05:00
ad025c1668 Merge pull request 'refactor: чистка мёртвого кода, лимит ходов D3, HTML-экранирование и UTF-8 обрезка в Telegram' (#4) from feat/83bf5b6fbb8b238c into main
Some checks failed
CI / test (push) Failing after 31s
CI / build-and-package (amd64, linux) (push) Successful in 35s
CI / build-and-package (amd64, windows) (push) Successful in 35s
Reviewed-on: http://gitea.hal9000.home/kamelion/ratatoskr-go/pulls/4
2026-08-18 23:14:50 +05:00
52 changed files with 3706 additions and 488 deletions

View File

@@ -1,10 +1,12 @@
# conventions # conventions
## Стиль / кодстайл ## Стиль / кодстайл
- Стандартный Go-стиль; гофм `gofmt`/`go fmt ./...`. Документация-комментарии и - Стандартный Go-стиль; `gofmt`/`go fmt ./...`. Документация-комментарии и
package-doc на русском языке (в START-комментариях файлов и doc-комментариях). package-doc на русском языке.
- Типизация: строгие типы, интерфейсы для абстракций (Decider, LiveProber). - Типизация: строгие типы, интерфейсы для абстракций (Decider, ui.Store, events.Event, chat.Channel).
- Свой тип `config.Duration` для времени (YAML-строки "5s"/"10m"), метод `.Duration()`. - Свой тип `config.Duration` для времени (YAML-строки "5s"/"10m"), метод `.Duration()`.
- UI: uni-directional data flow — UI читает снапшоты (копии, `ui.Store`), не мутирует storage/core.
Ядро публикует события в `events.Bus`, UI подписывается `events.Hub` → колбэки в UI-поток через `fyne.Do`.
## Обработка ошибок ## Обработка ошибок
- Ошибки классифицируются по идентификаторам классов в исходниках (см. ниже). - Ошибки классифицируются по идентификаторам классов в исходниках (см. ниже).
@@ -15,9 +17,10 @@
| Блок | Коды | Где | | Блок | Коды | Где |
|---|---|---| |---|---|---|
| C | C1C4 | internal/config | | C | C1C4 | internal/config |
| A | A1A4 | internal/analyst | | A | A1A4 | internal/analyst; A (app): A1A4 в internal/app |
| M | M1M5 | internal/chat | | M | M1M5 | internal/chat |
| O | O1O4 | internal/opencode | | D | D1D5 | internal/core (ErrDecideFailed, ErrMaxTurns, ErrCommandUnknown…) |
| O | O1O5 | internal/opencode |
| S | S1S5 | internal/storage | | S | S1S5 | internal/storage |
| W | W1W5 | internal/worker | | W | W1W5 | internal/worker |
| E | E1E4 | internal/worker (репозитории) | | E | E1E4 | internal/worker (репозитории) |
@@ -25,15 +28,16 @@
| U | U1U6 | internal/update | | U | U1U6 | internal/update |
## Архитектурные конвенции ## Архитектурные конвенции
- **Composition root** — internal/app; подсистемы собираются там, лимиты Core - **Composition root** — internal/app; подсистемы собираются там. UI-окно создаётся через
(MaxTurns=15, MaxConfirmCycles=3, MaxQuestionsPerTurn=5) в core.New. `newUIWindow(store)` (build-tag: cgo → desktop.New, !cgo → nil/headless). Флаг `--noui` выключает окно.
- **Агенты** (analyst.md/dev.md/reviewer.md) — markdown-промпты, встроены через go:embed - **Слои:** model (домен) ← storage/events/UI зависят только от model; core не знает про UI.
(internal/agents/*.md + embed.go), распаковываются в config_dir. Вердикт — строгий JSON. - **Агенты** (analyst.md/dev.md/reviewer.md) — markdown-промпты, go:embed, распаковка в config_dir.
Вердикт — строгий JSON.
- Репо-клонирование/воркеры — через gitops; ветки задач `feat/<taskTag>` (см. mem:core). - Репо-клонирование/воркеры — через gitops; ветки задач `feat/<taskTag>` (см. mem:core).
- Обратная совместимость: поле `Repo` (одиночный) и `Repos` (список); EffectiveRepos/ - Обратная совместимость: поле `Repo` (одиночный) и `Repos` (список); EffectiveRepos/
SetReposFromDB/ReposJoined в storage/models.go. SetReposFromDB/ReposJoined в storage/models.go.
## Версии ## Версии
- `app.Version` — семантическая major.minor.patch (ручной инкремент: patch=фиксы, - `app.Version` — семантическая major.minor.patch (ручной инкремент: patch=фиксы,
minor=новая обратно-совместимая функциональность, major=несовместимые изменения). minor=новая обратно-совместимая функциональность, major=несовместимые изменения). Сейчас 0.2.2.
- `main.version` (ldflag) — build-идентификатор `commit-<sha7>`, отдельно от app.Version. - `main.version` (ldflag) — build-идентификатор `commit-<sha7>`, отдельно от app.Version.

View File

@@ -1,51 +1,57 @@
# core # core
Ratatoskr-go — оркестратор конвейера Ratatoskr (порт с Python на Go) в единый Ratatoskr-go — оркестратор конвейера Ratatoskr (порт с Python на Go) в единый
статический бинарь (CGO_ENABLED=0). Субагенты запускаются через внешний процесс бинарь. Субагенты (analyst/dev/reviewer) запускаются через внешний процесс
[opencode](https://opencode.ai). Взаимодействие — Telegram-бот. [opencode](https://opencode.ai) (v2 HTTP API). Интерфейсы: Telegram-бот + Fyne UI.
Linux — статический headless-бинарь (CGO_ENABLED=0); Windows — с Fyne-окном (cgo).
## Структура (модули internal/) ## Структура (модули internal/)
``` ```
cmd/ratatoskr/ точка входа, сборка бинаря; main.version и main.updateToken вшиваются ldflag'ом cmd/ratatoskr/ точка входа; флаги: -config, -version, -noui (headless). main.version/updateToken вшиваются ldflag'ом
internal/ internal/
app/ composition root/DI: App.New -> config.Load+Validate, ResolveExePaths, storage, Runner, Analyst, Core, Worker, Updater. packageOwner="kamelion", Version="0.1.0" app/ composition root/DI: App.New -> config.Load+Validate, ResolveExePaths, storage, Runner, Analyst, Core, Worker, Updater, UI. packageOwner="kamelion", Version="0.2.2"
config/ YAML+env загрузка (${VAR:-default}), defaults, validate C1-C4 config/ YAML+env загрузка (${VAR:-default}), defaults, validate C1-C4
model/ доменные типы/контракты: Status + validTransitions (IsValidTransition/IsTerminal), TraceStatus. Нижний слой — не зависит от storage/events
chat/ мультиканальный Router; telegram — long-poll канал. Коды M1-M5 chat/ мультиканальный Router; telegram — long-poll канал. Коды M1-M5
core/ state-machine задач + Decider/analyst интерфейс. Коды D3/D4 core/ state-machine задач + Decider/analyst интерфейс. Коды D3/D4
analyst/ аналитик: промпт + разбор JSON-вердикта (opencode agent). Коды A1-A4 analyst/ аналитик: промпт + разбор JSON-вердикта (opencode agent). Коды A1-A4
events/ шина событий UI: Bus (buffered Sub/Pub), типизированный Hub, события TaskCreated/TaskStatusChanged/HistoryAppended/TraceAppended/AgentActivity, LogBus (панель «Логи»)
worker/ polling-планировщик + dev/reviewer-конвейер + gitops. Коды W*, E*, R* worker/ polling-планировщик + dev/reviewer-конвейер + gitops. Коды W*, E*, R*
agents/ встроенные opencode-агенты (analyst.md, dev.md, reviewer.md) через go:embed agents/ встроенные opencode-агенты (analyst.md, dev.md, reviewer.md) через go:embed
opencode/ обёртка запуска opencode, LiveRegistry, парсинг вердикта. Коды O1-O4 opencode/ HTTP-клиент v2 API opencode serve, поллинг вердикта, LiveRegistry. Коды O1-O5
storage/ SQLite (modernc.org/sqlite, без CGO): tasks, traces, task_history. Коды S1-S5 storage/ SQLite (modernc.org/sqlite, без CGO): tasks, traces, task_history. Коды S1-S5
update/ автообновление из Gitea Packages. Коды U1-U6 update/ автообновление из Gitea Packages. Коды U1-U6
ui/ desktop-оболочка: Store-абстракция (копии), Commands (UI->Core), Controller (шины->View), View/Window интерфейсы
ui/desktop/ Fyne-реализация окна (build-tag cgo, --noui / нет cgo → headless)
scripts/ build-publish-ui.ps1 — локальная Windows-сборка Fyne-бинаря + публикация в Gitea
docs/ ui-spec.md — спека Fyne UI (слои, event-bus, fyne.Do/снапшоты)
``` ```
## Ключевые инварианты ## Ключевые инварианты
- **Команды/ядро:** Core.ProcessTurn — state-machine поверх storage. Команды: - **App.New-сигнатура:** `App.New(configPath, version, updateToken string, noUI bool)` (4-й параметр — headless; cgo-вариант собирается только при наличии C-компилятора).
/start /cancel /skip /retry N /status N /continue N. Активная задача — одна на чат. - **UI:** окно — ещё одна реализация `chat.Channel` (присоединяется в Router). Core не трогает UI; обмен — событийная шина (events). Кнопка «Завершить» = полный выход (SetOnQuit→cancel→UI.Run возвращается); закрытие крестиком = сворачивание, Core живёт. UI собирается с `--noui`/без cgo.
- **Фазы аналитика (Decision.Phase):** `ask` (уточняющие вопросы), `propose` (правки - **Фазы аналитика (Decision.Phase):** `ask`, `propose`, `ready` (два последних обрабатываются одинаково в core), `abort`. Требования валидатора: ask — chat_reply/questions; propose — хотя бы одно изменённое поле; ready — без изменённых полей.
черновика), `ready` (черновик полон как есть), `abort` (тема вне проекта). propose и ready - **Статусы задач (internal/model):** draft→collecting→ready→approved→running→success/failed/timeout + cancelled/aborted/closed (терминальные). UserID — chat.ID (одна активная задача на чат).
в core обрабатываются одинаково. - **Decider/Worker/Analyst/Reviewer:** Decider=analyst интерфейс; Worker — polling-планировщик; Reviewer проверяет diff dev-ветки (R1-R6), вердикт JSON {passed, critical_issues, solid_violations, comments}.
- **Статусы задач:** draft→collecting→ready→approved→running→success/failed/timeout, - **gitops (worker):** worktree-режим; feature-ветка `feat/<taskTag>` от origin/main; push через http.extraHeader, токен Bearer.
плюс cancelled/aborted/closed. approved — финальное одобрение («создавай»), после чего - **Пути «всё рядом с .exe»:** db/worktree резолвятся от ExeDir; config.yaml — рядом с бинарём, фоллбэк cwd.
воркер берёт задачу. IsValidTransition/IsTerminal в storage/models.go. - **Автообновление:** авто = только Check+уведомление; замена — по /update; версии в `commit-<sha7>/` (не `latest/`); Verify сверяет предprod-версию (binary+в.в) .
- **Decider/Worker/Analyst/Reviewer:** Decider=analyst интерфейс (analyst пакет реализует);
Worker — polling-планировщик dev-агента; Reviewer проверяет diff dev-ветки (R1-R6), ## opencode (v2 HTTP API, >= 1.18.18)
вердикт JSON {passed, critical_issues, solid_violations, comments}.
- **gitops (worker):** worker работает с worktrees; feature-ветка = `feat/<taskTag>` от - Интеграция с субагентами — через headless `opencode serve`, **v2 API** (`/api/*`). Версия opencode >= 1.18.18.
origin/main; ensureBranch (reset --hard + clean -fd + checkout -B), branchDiff - **Хардпин модели:** при создании сессии читается top-level `model` из конфига opencode (`internal/opencode/config.go`, JSONC-стрип) и передаётся в `POST /api/session` как `{"model":{providerID,id}}`.
(`origin/main...<branch>`), push через http.extraHeader, токен как Bearer. - **О5 WARN (устойчивость к v1-конфигу):** конфиг по старой схеме молча игнорируется v2; провайдер без api → unsupported модели → fallback. Ratatoskr не чинит сам, но логирует warning; фактическая модель ответа сравнивается с ожидаемой. Правильный v2-вид: `api:{type:"aisdk",package,url}`, `request.headers` вместо `options.headers`.
- **Пути «всё рядом с .exe»:** относительные db/worktree резолвятся от каталога бинаря - **Поллинг вердикта:** `POST /api/session/:id/prompt` (durable admit) → `GET /api/session/:id/message?order=desc&limit=200` (новые assistant-сообщения, текст в `content[].type=="text"`) → завершение = `GET /api/session/active` без сессии + финальное assistant-сообщение, стабильное `settlePolls=2` опроса. `POST .../interrupt` вместо abort.
(ExeDir), не от cwd. config.yaml ищутся рядом с бинарём, фоллбэк cwd.
- **Автообновление:** авто = только Check+уведомление; замена — по /update. Версии в
Gitea Packages `commit-<sha7>/` (не `latest/`). Вердикты агентов — строгий JSON.
## Контракты (не ломать) ## Контракты (не ломать)
- `App.New(configPath, version, updateToken string)` — сигнатура. - `App.New(configPath, version, updateToken string, noUI bool)` — сигнатура.
- `packageOwner` — константа "kamelion" (не плодить vars/ldflag/конфиг). - `packageOwner` — константа "kamelion" (не плодить vars/ldflag/конфиг).
- `update.Updater` — создаётся структурой `&update.Updater{...}`, конструктора нет. - `update.Updater` — создаётся структурой `&update.Updater{...}`, конструктора нет.
- `ui.Store` — чтение-модель для UI: ListTasks/GetTask/GetHistory/GetTraces (копии, без ссылок на storage).
- `events.Bus`/`events.LogBus` — шина UI (см. internal/events). publisher-интерфейсы подключаются к Core/Worker.
См. также `mem:tech_stack`, `mem:conventions`, `mem:task_completion`, `mem:suggested_commands`. см. также `mem:tech_stack`, `mem:conventions`, `mem:task_completion`, `mem:suggested_commands`.

View File

@@ -6,22 +6,27 @@
- `go fmt ./...` (или `make fmt`). - `go fmt ./...` (или `make fmt`).
- `go build -o ratatoskr ./cmd/ratatoskr/` (или `make build`). - `go build -o ratatoskr ./cmd/ratatoskr/` (или `make build`).
- `./ratatoskr -config config.yaml` (или `make run`). - `./ratatoskr -config config.yaml` (или `make run`).
- `./ratatoskr -version` показать версию бинаря. - `./ratatoskr -version` — версия бинаря; `./ratatoskr -noui` — headless (без окна).
- Кросс-сборка: `make cross` (linux/amd64 + windows/amd64). - Кросс-сборка: `make cross` (linux/amd64 + windows/amd64, headless).
## Windows / UI (Fyne)
- Сборка с UI (cgo): `powershell -ExecutionPolicy Bypass -File scripts/build-publish-ui.ps1`.
Требует `scripts/.gitea-creds` (не коммитить) и MinGW/gcc в PATH. См. `mem:toolchain/cgo-winlibs-gcc`.
- Обычная headless-сборка на Windows без cgo: `go build -o ratatoskr.exe ./cmd/ratatoskr/`.
## Go-тулчейн ## Go-тулчейн
Хост без `go` в PATH — экспорт вручную: Хост без `go` в PATH — экспорт вручную:
``` ```
export PATH=/opt/data/.local/go/bin:$PATH export PATH=/opt/data/.local/go/bin:$PATH
``` ```
На машине разработки (Windows, cmd) — обычный system Go. На машине разработки (Windows) — обычный system Go.
## git (worktree-процесс Ratatoskr) ## git (worktree-процесс Ratatoskr)
- Feature-ветка: `feat/<taskTag>`, база — `origin/main`. Пример проверки diff всей ветки: - Feature-ветка: `feat/<taskTag>`, база — `origin/main`. Проверка diff всей ветки:
`git diff origin/main...HEAD`. `git diff origin/main...HEAD`.
- Проверить состав отслеживаемых файлов (например наличие .serena): `git ls-files`. - Проверить состав отслеживаемых файлов (например наличие .serena): `git ls-files`.
- Статус: `git status`. Лог: `git log --oneline -10`. - Статус: `git status`. Лог: `git log --oneline -10`.
## Замечания про среду ## Замечания про среду
- ОС Windows + PowerShell 5.1 (shell: powershell) — команды собирать/запускать с учётом - ОС Windows + PowerShell 5.1 (shell: powershell) — нет `&&`; использовать `;`/`if ($?) {}`;
этом (нет `&&`; использовать `;`/`if ($?) {}`; & для путей с пробелами). `&` для вызова путей с пробелами. Для серенных файлов лучше отдельная команда `serena memories check`.

View File

@@ -6,9 +6,10 @@
2. **Тесты:** `go test ./... -v -count=1 -timeout 120s` — все проходят (`make test`). 2. **Тесты:** `go test ./... -v -count=1 -timeout 120s` — все проходят (`make test`).
3. **Статический анализ:** `go vet ./...` — чисто (`make vet`). 3. **Статический анализ:** `go vet ./...` — чисто (`make vet`).
4. **Сборка:** `go build -o ratatoskr ./cmd/ratatoskr/` (`make build`) — компилируется. 4. **Сборка:** `go build -o ratatoskr ./cmd/ratatoskr/` (`make build`) — компилируется.
(Кросс-сборка `make cross` — только при необходимости.) - UI-сборка (cgo+Fyne) проверяется отдельно (`scripts/build-publish-ui.ps1` или `go build` с MinGW);
headless-сборка cgo-пакетов (`internal/ui/desktop`) — только typecheck.
5. **Пересмотр контрактов:** если менялась сигнатура `App.New(configPath, version, 5. **Пересмотр контрактов:** если менялась сигнатура `App.New(configPath, version,
updateToken string)` — обновить вызовы в `internal/app/app_test.go` (иначе go vet падает). updateToken string, noUI bool)` — обновить все вызовы (app_test.go, main.go).
6. Коммит осмысленными атомарными коммитами в feature-ветку `feat/<taskTag>` от origin/main. 6. Коммит осмысленными атомарными коммитами в feature-ветку `feat/<taskTag>` от origin/main.
Рабочий процесс Ratatoskr (агент dev в этом конвейере): изучить код, реализовать так, чтобы Рабочий процесс Ratatoskr (агент dev в этом конвейере): изучить код, реализовать так, чтобы

View File

@@ -1,30 +1,32 @@
# tech_stack # tech_stack
## Язык / рантайм ## Язык / рантайм
- Go **1.25.0** (go.mod `go 1.25.0`). Модуль `github.com/kamelion/ratatoskr-go`. - Go **1.25.0** (go.mod `go 1.25.0`, module `github.com/kamelion/ratatoskr-go`).
- Сборка: статический бинарь, `CGO_ENABLED=0`. Локальный Go-тулчейн на хосте: - Headless-сборка статический бинарь `CGO_ENABLED=0`. Windows-сборка с UI — cgo (Fyne/GLFW) + MinGW-w64/gcc.
`export PATH=/opt/data/.local/go/bin:$PATH` (go1.25.0 linux/amd64). - Локальный Go-тулчейн на хосте: `export PATH=/opt/data/.local/go/bin:$PATH` (go1.25.0 linux/amd64). На Windows — обычный system Go.
## Основные зависимости (go.mod) ## Основные зависимости (go.mod)
- `fyne.io/fyne/v2 v2.6.0` — desktop UI (cgo/GLFW; платформа Windows + Linux). Не собирается в CI (headless).
- `gopkg.in/yaml.v3 v3.0.1` — парсинг config.yaml. - `gopkg.in/yaml.v3 v3.0.1` — парсинг config.yaml.
- `modernc.org/sqlite v1.56.0` — SQLite без CGO (чистый Go). - `modernc.org/sqlite v1.56.0` — SQLite без CGO (чистый Go).
- (indirect) google/uuid, go-humanize, mattn/go-isatty, x/sys, modernc.org/libc/mathutil/memory. - (indirect) google/uuid, go-humanize, go-text/typesetting, x/sys, x/net, modernc.org/libc и пр.
## Внешние процессы ## Внешние процессы
- **opencode** (opencode.ai) — внешний процесс для субагентов (analyst/dev/reviewer). - **opencode** (opencode.ai) — субагенты (analyst/dev/reviewer). Управляется через internal/opencode (Runner, Pool, LiveRegistry). Требуемая версия **>= 1.18.18** (v2 HTTP API `/api/*`).
Управляется через internal/opencode (Runner, LiveRegistry). Настраивается в конфиге - Настройки в конфиге: `opencode.bin / config / config_dir / hard_timeout / idle_timeout / poll_ms` + `opencode.server` (hostname, port, password).
(opencode.bin / config / config_dir / hard_timeout / idle_timeout / poll_ms). - Подробно — `mem:core`.
## Сборка / Makefile ## Сборка / Makefile
- `make build` — go build -ldflags="-s -w -X main.version=commit-<sha7> -X main.updateToken=..." -o ratatoskr ./cmd/ratatoskr/. - `make build` — go build `-ldflags="-s -w -X main.version=commit-<sha7> -X main.updateToken=..." -o ratatoskr ./cmd/ratatoskr/`.
- `make test` — go test ./... -v -count=1 -timeout 120s. - `make test``go test ./... -v -count=1 -timeout 120s`.
- `make vet` — go vet ./... `make fmt` — go fmt ./... - `make vet``go vet ./...`; `make fmt``go fmt ./...`.
- `make run` — build + ./ratatoskr -config config.yaml. - `make run` — build + `./ratatoskr -config config.yaml`.
- `make cross` кросс-сборка linux/amd64 + windows/amd64 (ratatoskr-windows-amd64.exe — Windows-машина Камиля). - `make cross` — linux/amd64 + windows/amd64 (headless, БЕЗ Fyne-UI; для UI нужен cgo).
- GIT_SHA вшивается в main.version; UPDATE_TOKEN — в main.updateToken (секрет только у CI). - **Fyne-UI (Windows) собирается вручную** на машине с MinGW:
`scripts/build-publish-ui.ps1` (cgo) — сборка + публикация в Gitea в версию `commit-<sha7>`.
Требует `scripts/.gitea-creds` (git-ignored) и gcc в PATH. См. `mem:toolchain/cgo-winlibs-gcc`.
## CI (.gitea/workflows/ci.yaml) ## CI (.gitea/workflows/ci.yaml)
- Job `test`: go test ./... + go vet ./.... - Job `test`: `go test ./...` + `go vet ./...` (setup-go '1.23', check-latest).
- Job `build-and-package` (matrix linux/amd64+windows/amd64): собирает и публикует в - Job `build-and-package` (matrix linux/amd64 + windows/amd64): собирает headless-бинарь с ldflag update_token, публикует на `main` в версию `commit-<sha7>/` (+ companion `.version`/`.sha256`).
Gitea Packages на `main` в версию `commit-<sha7>/` + companion-файлы `.version`/`.sha256`. - Секреты: `TC_GITEA_TOKEN` (write:packages), `TC_UPDATE_TOKEN` (read:package), vars `GIT_MAIN_URL`.
- Секреты: TC_GITEA_TOKEN (write:packages), TC_UPDATE_TOKEN (read:package), GIT_MAIN_URL.

View File

@@ -0,0 +1,13 @@
# cgo/winlibs-gcc
## MinGW/WinLibs для Fyne-сборки на Windows
- UI (internal/ui/desktop, fyne.io/fyne/v2) требует cgo + C-компилятор для GLFW.
- Используется **WinLibs** (mingw64) gcc. Bin каталог нужно добавить в PATH перед
`scripts/build-publish-ui.ps1`:
```powershell
$p = "<...>\mingw64\bin"; $env:Path = "$p;$env:Path"
```
- Проверка: `gcc --version`; также убедиться, что `CGO_ENABLED` не = "0".
- Care: не использовать Cygwin/MSYS-сборки gcc — только native MinGW-w64 (WinLibs).
- Сборка UI выполняется локально на Windows-машине, публикуется в Gitea Packages
в в версию `commit-<sha7>` со companion-.version/.sha256 (контракт U4) — см. скрипт и `mem:tech_stack`.

View File

@@ -1,39 +1,6 @@
# the name by which the project can be referenced within Serena # the name by which the project can be referenced within Serena/when chatting with the LLM.
project_name: "ratatoskr-go" project_name: "ratatoskr-go"
# list of languages for which language servers are started; choose from:
# al angular ansible bash clojure
# cpp cpp_ccls crystal csharp csharp_omnisharp
# dart elixir elm erlang fortran
# fsharp go groovy haskell haxe
# hlsl html java json julia
# kotlin lean4 lua luau markdown
# matlab msl nix ocaml pascal
# perl php php_phpactor powershell python
# python_jedi python_ty r rego ruby
# ruby_solargraph rust scala scss solidity
# svelte swift systemverilog terraform toml
# typescript typescript_vts vue yaml zig
# (This list may be outdated. For the current list, see values of Language enum here:
# https://github.com/oraios/serena/blob/main/src/solidlsp/ls_config.py
# For some languages, there are alternative language servers, e.g. csharp_omnisharp, ruby_solargraph.)
# Note:
# - For C, use cpp
# - For JavaScript, use typescript
# - For Angular projects, use angular (subsumes typescript+html; requires `npm install` in the project root)
# - For Svelte projects, use svelte (subsumes typescript/javascript for .svelte projects; requires npm)
# - For SCSS / Sass / plain CSS, use scss (some-sass-language-server handles all three)
# - For Free Pascal/Lazarus, use pascal
# Special requirements:
# Some languages require additional setup/installations.
# See here for details: https://oraios.github.io/serena/01-about/020_programming-languages.html#language-servers
# When using multiple languages, the first language server that supports a given file will be used for that file.
# The first language is the default language and the respective language server will be used as a fallback.
# Note that when using the JetBrains backend, language servers are not used and this list is correspondingly ignored.
languages:
- go
# the encoding used by text files in the project # the encoding used by text files in the project
# For a list of possible encodings, see https://docs.python.org/3.11/library/codecs.html#standard-encodings # For a list of possible encodings, see https://docs.python.org/3.11/library/codecs.html#standard-encodings
encoding: "utf-8" encoding: "utf-8"
@@ -55,23 +22,19 @@ ignore_all_files_in_gitignore: true
# advanced configuration option allowing to configure language server-specific options. # advanced configuration option allowing to configure language server-specific options.
# Maps the language key to the options. # Maps the language key to the options.
# Have a look at the docstring of the constructors of the LS implementations within solidlsp (e.g., for C# or PHP) to see which options are available. # The settings are considered only if the project is trusted (see global configuration to define trusted projects).
# No documentation on options means no options are available. # See https://oraios.github.io/serena/02-usage/050_configuration.html#language-server-specific-settings
ls_specific_settings: {} ls_specific_settings: {}
# list of additional workspace folder paths for cross-package reference support (e.g. in monorepos).
# Paths can be absolute or relative to the project root.
# Each folder is registered as an LSP workspace folder, enabling language servers to discover
# symbols and references across package boundaries.
# Currently supported for: TypeScript.
# Example:
# additional_workspace_folders:
# - ../sibling-package
# - ../shared-lib
additional_workspace_folders: []
# list of additional paths to ignore in this project. # list of additional paths to ignore in this project.
# Same syntax as gitignore, so you can use * and **. # Same syntax as gitignore, so you can use * and **.
# Important: quote patterns that start with `*`, otherwise YAML treats them as aliases.
# Example:
# ignored_paths:
# - "examples/**"
# - ".worktrees/**"
# - "**/bin/**"
# - "**/obj/**"
# Note: global ignored_paths from serena_config.yml are also applied additively. # Note: global ignored_paths from serena_config.yml are also applied additively.
ignored_paths: [] ignored_paths: []
@@ -131,3 +94,76 @@ read_only_memory_patterns: []
# Extends the list from the global configuration, merging the two lists. # Extends the list from the global configuration, merging the two lists.
# Example: ["_archive/.*", "_episodes/.*"] # Example: ["_archive/.*", "_episodes/.*"]
ignored_memory_patterns: [] ignored_memory_patterns: []
# list of additional workspace folder paths for cross-package reference support.
# Paths can be absolute or relative to the project root.
# Each folder is registered as an LSP workspace folder, enabling language servers to discover
# symbols and references across package boundaries, but these folders are not indexed by Serena,
# i.e. the respective symbols will not be found using Serena's symbol search tools.
# Example:
# additional_workspace_folders:
# - ../sibling-package
# - ../shared-lib
ls_additional_workspace_folders: []
# list of language servers to start when using the LSP backend; choose from:
# ada al angular ansible bash
# bsl clojure cpp cpp_ccls crystal
# csharp csharp_omnisharp cue dart deno
# elixir elm erlang fortran fsharp
# gdscript gleam go groovy haskell
# haxe hlsl html java json
# julia kotlin latex lean4 lua
# luau markdown matlab msl nextflow
# nix ocaml pascal perl php
# php_phpactor php_phpantom powershell python python_basedpyright
# python_jedi python_pyrefly python_ty qml r
# rego ruby ruby_solargraph rust scala
# scss solidity svelte swift systemverilog
# terraform toml typescript typescript_vts vue
# wolfram yaml zig
# (This list may be outdated; generated with scripts/print_language_list.py;
# For the current list, see values of the LanguageServerId enum here:
# https://github.com/oraios/serena/blob/main/src/solidlsp/ls_config.py)
# For some languages, there are several alternative language servers, e.g. csharp_omnisharp, ruby_solargraph.)
# Note:
# - For C, use cpp
# - For JavaScript, use typescript
# - For Angular projects, use angular (subsumes typescript+html; requires `npm install` in the project root)
# - For Svelte projects, use svelte (subsumes typescript/javascript for .svelte projects; requires npm)
# - For Deno projects, use deno (serves the same .ts/.js files as typescript; requires the deno CLI on PATH)
# - For SCSS / Sass / plain CSS, use scss (some-sass-language-server handles all three)
# - For Free Pascal/Lazarus, use pascal
# Special requirements:
# Some language servers require additional setup/installations.
# See here for details: https://oraios.github.io/serena/01-about/020_programming-languages.html#language-servers
# When using multiple language servers, the first language server that supports a given file will be used for that file.
# The first language server is the default language and the respective language server will be used as a fallback.
# Note that when using the JetBrains backend, language servers are not used and this list is correspondingly ignored.
language_servers:
- go
# list of workspace folder paths (LSP backend only).
# These folders will be used to build up Serena's symbol index.
# Paths must be within the project root and should thus be relative to the project root.
# Furthermore, the paths should not be filtered by ignore settings.
# Default setting: The entire project root folder (".") is considered.
# In (large) monorepos, this can be used to index only subfolders of the project root, e.g.
# ls_workspace_folders:
# - "./subproject1"
# - "./subproject2"
ls_workspace_folders:
- .
# optional shell command to run before the language backend (LSP or JetBrains) is initialised.
# the command runs in the project root directory and is only executed if the project is trusted
# (see trusted_project_path_patterns in the global configuration).
# serena waits for the command to exit: a non-zero exit code is logged as an error but does not
# abort activation. a per-project timeout (activation_command_timeout, default 180s) is the safety
# backstop for non-terminating commands; on expiry the process is killed and activation continues.
# example: activation_command: "npx nx run-many -t build"
activation_command:
# maximum time in seconds to wait for activation_command to complete before killing it (default 180s).
# must be a positive number.
activation_command_timeout: 180.0

View File

@@ -17,7 +17,7 @@ internal/
analyst/ # аналитик: промпт + разбор JSON-решения (opencode agent, классы A1A4) analyst/ # аналитик: промпт + разбор JSON-решения (opencode agent, классы A1A4)
worker/ # polling-планировщик + dev/reviewer-конвейер (классы W*, E*, R*) worker/ # polling-планировщик + dev/reviewer-конвейер (классы W*, E*, R*)
agents/ # встроенные агенты opencode (analyst.md, dev.md, reviewer.md) через go:embed agents/ # встроенные агенты opencode (analyst.md, dev.md, reviewer.md) через go:embed
opencode/ # обёртка запуска opencode-процесса, парсинг вердикта (классы O1O4) opencode/ # HTTP-клиент v2 API opencode serve, поллинг вердикта (классы O1O5)
storage/ # SQLite (modernc.org/sqlite, без CGO): tasks, traces, task_history (S1S5) storage/ # SQLite (modernc.org/sqlite, без CGO): tasks, traces, task_history (S1S5)
update/ # автообновление из Gitea Packages (классы U1U6) update/ # автообновление из Gitea Packages (классы U1U6)
``` ```
@@ -129,6 +129,33 @@ update:
| `/status` | версия бинаря + есть ли доступное обновление | | `/status` | версия бинаря + есть ли доступное обновление |
| `/help` | справка по всем командам | | `/help` | справка по всем командам |
## Интеграция с opencode (субагенты)
Субагенты (analyst / dev / reviewer) запускаются через **headless** `opencode serve`
по **v2 HTTP API** (префикс `/api/*`). Требуемая версия opencode: **>= 1.18.18**
(сборки с v2 HTTP API). Старый бинарь, отвечающий только на `/global/health`,
не подходит: healthcheck падает с понятной ошибкой (класс O1).
Что делает обёртка (`internal/opencode`):
- **Хардпин модели.** При создании сессии в конфиге opencode ищется top-level
`"model"` (`internal/opencode/config.go`) и передаётся в `POST /api/session`
как `{"model":{providerID,id}}`. Это убирает зависимость от fallback-логики
opencode (которая молча выбирает «дефолтную» запись, если модель не задана).
- **Весь код резолва модели устойчив к этому классу проблем (класс O5 WARN):**
- если конфиг не читается / в нём нет `model` — в логи пишется warning;
- фактическая модель ответа (из финального assistant-сообщения) сравнивается
с ожидаемой; расхождение логируется как warning;
- конфиг, написанный по **старой v1-схеме** (`provider.X.npm` / `options`),
молча игнорируется v2 — обёртка этого не «чинит» сама, но предупреждает.
Правильный v2-вид провайдера — `api: { type:"aisdk", package, url }` и
`request.headers` вместо `options.headers`.
- **Поллинг вердикта.** Промпт отправляется неблокирующе (`POST .../prompt`
durable admit), вердикт собирается из новых assistant-сообщений
(`GET .../message`); завершение ответа — сессия ушла из активных дренажей
(`GET .../active`) и появилось финальное assistant-сообщение, стабильное
несколько опросов подряд.
## Фазы аналитика ## Фазы аналитика
Аналитик (`internal/analyst`) возвращает JSON-вердикт с полем `phase`: Аналитик (`internal/analyst`) возвращает JSON-вердикт с полем `phase`:
@@ -198,7 +225,7 @@ curl -s -H "Authorization: token $TOKEN" "$B/api/packages/kamelion/generic/ratat
| C | C1C4 | `internal/config` | | C | C1C4 | `internal/config` |
| A | A1A4 | `internal/analyst` | | A | A1A4 | `internal/analyst` |
| M | M1M5 | `internal/chat` | | M | M1M5 | `internal/chat` |
| O | O1O4 | `internal/opencode` | | O | O1O5 | `internal/opencode` |
| S | S1S5 | `internal/storage` | | S | S1S5 | `internal/storage` |
| W | W1W5 | `internal/worker` | | W | W1W5 | `internal/worker` |
| E | E1E4 | `internal/worker` (репозитории) | | E | E1E4 | `internal/worker` (репозитории) |

View File

@@ -27,6 +27,7 @@ var updateToken = ""
func main() { func main() {
cfg := flag.String("config", "", "путь к config.yaml (по умолчанию — CWD/config.yaml)") cfg := flag.String("config", "", "путь к config.yaml (по умолчанию — CWD/config.yaml)")
versionFlag := flag.Bool("version", false, "показать версию и выйти") versionFlag := flag.Bool("version", false, "показать версию и выйти")
noUI := flag.Bool("noui", false, "работать без графического окна (headless)")
flag.Parse() flag.Parse()
if *versionFlag { if *versionFlag {
@@ -34,7 +35,7 @@ func main() {
return return
} }
a, err := app.New(*cfg, version, updateToken) a, err := app.New(*cfg, version, updateToken, *noUI)
if err != nil { if err != nil {
log.Fatalf("app init: %v", err) log.Fatalf("app init: %v", err)
} }

167
docs/ui-spec.md Normal file
View File

@@ -0,0 +1,167 @@
# 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. Скоуп реализации (поэтапно)
- [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.
---
## Открытые пункты (TODO)
- [ ] Точно определить set сплит-панелей (2x2 центр) и как добавляются
- [ ] Обязательные поля формы создания задачи
- [ ] Подробности виджета списка задач (после первого мильстон)
- [x] Архитектура обмена UI<->Core (раздел 12 — принципы; детальные контракты событий/интерфейсов — следующая итерация)

31
go.mod
View File

@@ -3,17 +3,48 @@ module github.com/kamelion/ratatoskr-go
go 1.25.0 go 1.25.0
require ( require (
fyne.io/fyne/v2 v2.6.0
gopkg.in/yaml.v3 v3.0.1 gopkg.in/yaml.v3 v3.0.1
modernc.org/sqlite v1.56.0 modernc.org/sqlite v1.56.0
) )
require ( 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/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/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/mattn/go-isatty v0.0.24 // indirect
github.com/ncruces/go-strftime v1.0.0 // 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/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/sys v0.47.0 // indirect
golang.org/x/text v0.22.0 // indirect
modernc.org/libc v1.74.4 // indirect modernc.org/libc v1.74.4 // indirect
modernc.org/mathutil v1.7.1 // indirect modernc.org/mathutil v1.7.1 // indirect
modernc.org/memory v1.11.0 // indirect modernc.org/memory v1.11.0 // indirect

74
go.sum
View File

@@ -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 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY=
github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= 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 h1:LMLX+LgTNWpfvCBdFebv6EsYotImrt/Ppc5cXIriCSo=
github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3/go.mod h1:jl5iWTm0/hd5PjEYEOuwAJ57L/CibdZfrqZ5XA5GrCk= 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 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= 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 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k=
github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= 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 h1:tGZZoVgT/KiqK1c8ocVLeDS8BSWMRd47J3Lbz7vsReI=
github.com/mattn/go-isatty v0.0.24/go.mod h1:nMCL3Zebbrt45jsMDgnfIwz6ydEQApk5oEI3HqDio6A= 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 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w=
github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls= 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 h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE=
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo= 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 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ=
golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0= 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 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM=
golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= 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 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= 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 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q=
golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA= 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 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 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
modernc.org/cc/v4 v4.29.1 h1:MKgdCV3WykTSPqpVrnxdEDS0HEd2FHpKZDzxzU5LyeI= modernc.org/cc/v4 v4.29.1 h1:MKgdCV3WykTSPqpVrnxdEDS0HEd2FHpKZDzxzU5LyeI=

View File

@@ -8,6 +8,7 @@ import (
"strings" "strings"
"github.com/kamelion/ratatoskr-go/internal/core" "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/opencode"
"github.com/kamelion/ratatoskr-go/internal/storage" "github.com/kamelion/ratatoskr-go/internal/storage"
) )
@@ -35,6 +36,16 @@ type Analyst struct {
Runner OpenCodeRunner Runner OpenCodeRunner
Worktree string // каталог, откуда запускать opencode run Worktree string // каталог, откуда запускать opencode run
Agent string // имя агента (default "analyst") 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-ответа аналитика. // AnalystResponse — структура JSON-ответа аналитика.
@@ -58,6 +69,8 @@ func (a *Analyst) Decide(ctx context.Context, history []core.Message, draft stor
agent = "analyst" agent = "analyst"
} }
a.publish(events.AgentActivity{TaskID: draft.ID, Agent: agent, Stage: "decide"})
if len(history) == 0 && !force { if len(history) == 0 && !force {
return core.Decision{}, fmt.Errorf("%w: пустая история диалога", ErrNotReady) return core.Decision{}, fmt.Errorf("%w: пустая история диалога", ErrNotReady)
} }

View File

@@ -5,6 +5,7 @@ import (
"context" "context"
"errors" "errors"
"fmt" "fmt"
"io"
"log" "log"
"os" "os"
"os/signal" "os/signal"
@@ -19,8 +20,10 @@ import (
"github.com/kamelion/ratatoskr-go/internal/chat/telegram" "github.com/kamelion/ratatoskr-go/internal/chat/telegram"
"github.com/kamelion/ratatoskr-go/internal/config" "github.com/kamelion/ratatoskr-go/internal/config"
"github.com/kamelion/ratatoskr-go/internal/core" "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/opencode"
"github.com/kamelion/ratatoskr-go/internal/storage" "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/update"
"github.com/kamelion/ratatoskr-go/internal/worker" "github.com/kamelion/ratatoskr-go/internal/worker"
) )
@@ -61,6 +64,15 @@ type App struct {
Updater *update.Updater Updater *update.Updater
tg *telegram.Channel // сохранена для Run tg *telegram.Channel // сохранена для Run
pool *opencode.Pool // пул opencode serve-серверов (API-режим) pool *opencode.Pool // пул opencode serve-серверов (API-режим)
// Events — доменная шина UI; LogEvents — шина логов (панель «Логи»).
Events *events.Bus
LogEvents *events.LogBus
// UI — десктопное окно (если собрано с cgo и не задан --noui).
// nil — headless-режим (без окна).
UI ui.Window
Controller *ui.Controller
} }
// New читает конфиг и собирает все зависимости. // New читает конфиг и собирает все зависимости.
@@ -68,7 +80,7 @@ type App struct {
// version — вшитая версия бинаря (ldflag -X main.version). // version — вшитая версия бинаря (ldflag -X main.version).
// updateToken — вшитый токен read:package для авто-обновления // updateToken — вшитый токен read:package для авто-обновления
// (ldflag -X main.updateToken); имеет приоритет над update.token из конфига. // (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) cfg, err := config.Load(configPath)
if err != nil { if err != nil {
return nil, fmt.Errorf("%w: %v", ErrConfig, err) return nil, fmt.Errorf("%w: %v", ErrConfig, err)
@@ -144,8 +156,37 @@ func New(configPath, version, updateToken string) (*App, error) {
CoreCtx: coreCtx, CoreCtx: coreCtx,
pool: ocPool, 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) 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-канал // Telegram-канал
tg := telegram.New(cfg.Telegram.Token, cfg.Chat.PollInterval.Duration()) tg := telegram.New(cfg.Telegram.Token, cfg.Chat.PollInterval.Duration())
a.tg = tg a.tg = tg
@@ -166,6 +207,7 @@ func New(configPath, version, updateToken string) (*App, error) {
GitToken: cfg.Git.Token, GitToken: cfg.Git.Token,
Live: live, Live: live,
Notify: a, // авто-уведомления владельцу задачи через Router Notify: a, // авто-уведомления владельцу задачи через Router
Events: a.Events,
} }
a.Router = router a.Router = router
a.Worker = w a.Worker = w
@@ -226,6 +268,14 @@ func (a *App) Run(ctx context.Context) error {
// Авто-проверка обновления (только уведомление владельца; замена — по /update) // Авто-проверка обновления (только уведомление владельца; замена — по /update)
a.startAutoCheck(ctx) 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 // Ожидание сигнала или фатальной ошибки Telegram
sigCh := make(chan os.Signal, 1) sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM) signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)

View File

@@ -21,10 +21,11 @@ func TestNew(t *testing.T) {
t.Setenv("RATATOSKR_DB", dbPath) t.Setenv("RATATOSKR_DB", dbPath)
// Загружаем без config-файла (дефолты + env) // Загружаем без config-файла (дефолты + env)
a, err := New("", "dev", "") a, err := New("", "dev", "", true)
if err != nil { if err != nil {
t.Fatalf("New() err = %v", err) t.Fatalf("New() err = %v", err)
} }
defer a.Store.Close()
if a.Store == nil { if a.Store == nil {
t.Fatal("Store не создан") t.Fatal("Store не создан")
} }
@@ -51,7 +52,7 @@ func TestNew_MissingToken(t *testing.T) {
t.Setenv("TG_CHAT_ID", "12345") t.Setenv("TG_CHAT_ID", "12345")
t.Setenv("RATATOSKR_DB", dbPath) t.Setenv("RATATOSKR_DB", dbPath)
_, err := New("", "dev", "") _, err := New("", "dev", "", true)
if err == nil { if err == nil {
t.Fatal("expected error for missing token") t.Fatal("expected error for missing token")
} }
@@ -66,7 +67,7 @@ func TestNew_BadDB(t *testing.T) {
t.Setenv("TG_CHAT_ID", "12345") t.Setenv("TG_CHAT_ID", "12345")
t.Setenv("RATATOSKR_DB", dbPath) t.Setenv("RATATOSKR_DB", dbPath)
_, err := New("", "dev", "") _, err := New("", "dev", "", true)
if err == nil { if err == nil {
t.Fatal("expected error for invalid db path") t.Fatal("expected error for invalid db path")
} }
@@ -83,10 +84,11 @@ func TestNew_RunCtxCancel(t *testing.T) {
t.Setenv("TG_CHAT_ID", "12345") t.Setenv("TG_CHAT_ID", "12345")
t.Setenv("RATATOSKR_DB", dbPath) t.Setenv("RATATOSKR_DB", dbPath)
a, err := New("", "dev", "") a, err := New("", "dev", "", true)
if err != nil { if err != nil {
t.Fatalf("New() err = %v", err) t.Fatalf("New() err = %v", err)
} }
defer a.Store.Close()
ctx, cancel := context.WithCancel(context.Background()) ctx, cancel := context.WithCancel(context.Background())
cancel() // сразу отменяем cancel() // сразу отменяем
@@ -120,15 +122,15 @@ func TestNew_WorktreeCreated(t *testing.T) {
" token: \"test:token\"", " token: \"test:token\"",
" chat_id: \"12345\"", " chat_id: \"12345\"",
"paths:", "paths:",
" db: \"" + filepath.Join(tmp, "test.db") + "\"", " db: '" + filepath.Join(tmp, "test.db") + "'",
" worktree: \"" + wt + "\"", " worktree: '" + wt + "'",
"", "",
}, "\n") }, "\n")
if err := os.WriteFile(configPath, []byte(content), 0o600); err != nil { if err := os.WriteFile(configPath, []byte(content), 0o600); err != nil {
t.Fatalf("write config: %v", err) t.Fatalf("write config: %v", err)
} }
a, err := New(configPath, "dev", "") a, err := New(configPath, "dev", "", true)
if err != nil { if err != nil {
t.Fatalf("New() err = %v", err) t.Fatalf("New() err = %v", err)
} }
@@ -165,7 +167,7 @@ func TestNew_UpdateWiring(t *testing.T) {
" base_url: \"https://hub.example.com\"", " base_url: \"https://hub.example.com\"",
" token: \"cfg-update-token\"", " token: \"cfg-update-token\"",
"paths:", "paths:",
" db: \"" + dbPath + "\"", " db: '" + dbPath + "'",
"", // пустая строка в конце "", // пустая строка в конце
}, "\n") }, "\n")
if err := os.WriteFile(configPath, []byte(content), 0o600); err != nil { if err := os.WriteFile(configPath, []byte(content), 0o600); err != nil {
@@ -173,10 +175,11 @@ func TestNew_UpdateWiring(t *testing.T) {
} }
// 1) base_url из update-блока, а НЕ git // 1) base_url из update-блока, а НЕ git
a, err := New(configPath, "dev", "") a, err := New(configPath, "dev", "", true)
if err != nil { if err != nil {
t.Fatalf("New() err = %v", err) t.Fatalf("New() err = %v", err)
} }
defer a.Store.Close()
if a.Updater == nil { if a.Updater == nil {
t.Fatal("Updater не создан") t.Fatal("Updater не создан")
} }
@@ -193,10 +196,11 @@ func TestNew_UpdateWiring(t *testing.T) {
} }
// 2) вшитый updateToken перекрывает конфиг // 2) вшитый updateToken перекрывает конфиг
a2, err := New(configPath, "dev", "embedded-update-token") a2, err := New(configPath, "dev", "embedded-update-token", true)
if err != nil { if err != nil {
t.Fatalf("New() err = %v", err) t.Fatalf("New() err = %v", err)
} }
defer a2.Store.Close()
if a2.Updater.Token != "embedded-update-token" { if a2.Updater.Token != "embedded-update-token" {
t.Errorf("Token = %q, want embedded-update-token (вшитый приоритетнее)", a2.Updater.Token) t.Errorf("Token = %q, want embedded-update-token (вшитый приоритетнее)", a2.Updater.Token)
} }

View File

@@ -36,22 +36,35 @@ import (
) )
// вердикты фейкового агента по имени. // вердикты фейкового агента по имени.
var ( var e2eAgentVerdicts = map[string]string{
e2eAgentVerdicts = map[string]string{
"analyst": `{"phase":"propose","title":"Калькулятор","goal":"Сделать веб-калькулятор","repo":"calc","why":"Нужен для учёта","ac":"Работает + - * /","chat_reply":"Черновик готов."}`, "analyst": `{"phase":"propose","title":"Калькулятор","goal":"Сделать веб-калькулятор","repo":"calc","why":"Нужен для учёта","ac":"Работает + - * /","chat_reply":"Черновик готов."}`,
"dev": `done`, "dev": `done`,
"reviewer": `{"passed":true,"comments":[]}`, "reviewer": `{"passed":true,"comments":[]}`,
} }
)
// e2eFakeAPI поднимает фейковый opencode serve experimental HTTP API (пути // e2eFakeAPI поднимает фейковый opencode serve, эмулирующий v2 HTTP API
// БЕЗ /api) и возвращает URL. По title сессии (ratatoskr-<agent>) определяет // (пути с префиксом /api/*, см. Client в internal/opencode). Агент
// агента и возвращает его вердикт как text-часть ответа на POST /message. // (analyst/dev/reviewer) определяется по тексту промпта на POST
// /api/session/{id}/prompt; вердикт возвращается как text-часть завершённого
// assistant-сообщения, которое отдаёт GET /api/session/{id}/message.
func e2eFakeAPI(t *testing.T) string { func e2eFakeAPI(t *testing.T) string {
t.Helper() t.Helper()
var mu sync.Mutex var mu sync.Mutex
sessions := map[string]string{} // id → agent sessions := map[string]string{} // id → agent
agentOf := func(prompt string) string {
switch {
case strings.Contains(prompt, "Ты — аналитик"):
return "analyst"
case strings.Contains(prompt, "Ты — dev-агент"):
return "dev"
case strings.Contains(prompt, "Ты — ревьюер"):
return "reviewer"
default:
return "unknown"
}
}
verdictFor := func(agent string) string { verdictFor := func(agent string) string {
if v, ok := e2eAgentVerdicts[agent]; ok { if v, ok := e2eAgentVerdicts[agent]; ok {
return v return v
@@ -59,39 +72,59 @@ func e2eFakeAPI(t *testing.T) string {
return "unknown agent" return "unknown agent"
} }
sessionID := func(path, suffix string) string {
return strings.TrimSuffix(strings.TrimPrefix(path, "/api/session/"), suffix)
}
assistantMsg := func(id, agent string) map[string]any {
ts := time.Now().UnixMilli()
return map[string]any{
"id": "m-" + id,
"type": "assistant",
"content": []map[string]any{{"type": "text", "text": verdictFor(agent)}},
"model": map[string]any{"providerID": "test", "id": "m"},
"time": map[string]any{"created": ts, "completed": ts},
}
}
h := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { h := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch { switch {
case r.Method == http.MethodPost && r.URL.Path == "/session": case r.Method == http.MethodPost && r.URL.Path == "/api/session":
// v2 create: {model:{...}} → {data:{id}}
id := fmt.Sprintf("e2e-%d", len(sessions)+1)
mu.Lock()
sessions[id] = ""
mu.Unlock()
writeJSON(w, map[string]any{"data": map[string]any{"id": id}})
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/prompt"):
// v2 durable admit: {prompt:{text}} → {data:{id,timeCreated}}
id := sessionID(r.URL.Path, "/prompt")
var req struct { var req struct {
Title string `json:"title"` Prompt struct {
Text string `json:"text"`
} `json:"prompt"`
} }
_ = json.NewDecoder(r.Body).Decode(&req) _ = json.NewDecoder(r.Body).Decode(&req)
agent := strings.TrimPrefix(req.Title, "ratatoskr-")
mu.Lock() mu.Lock()
id := fmt.Sprintf("e2e-%d", len(sessions)+1) sessions[id] = agentOf(req.Prompt.Text)
sessions[id] = agent
mu.Unlock() mu.Unlock()
// experimental: голая Session, id напрямую. writeJSON(w, map[string]any{"data": map[string]any{"id": "p-" + id, "timeCreated": time.Now().UnixMilli()}})
writeJSON(w, map[string]any{"id": id, "agent": agent, "model": map[string]any{"id": "m"}})
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/message"):
// блокирующий ответ: вердикт как text-часть.
id := strings.TrimSuffix(strings.TrimPrefix(r.URL.Path, "/session/"), "/message")
mu.Lock()
agent := sessions[id]
mu.Unlock()
writeJSON(w, map[string]any{"info": map[string]any{"role": "assistant"}, "parts": []map[string]any{{"type": "text", "text": verdictFor(agent)}}})
case r.Method == http.MethodGet && strings.HasSuffix(r.URL.Path, "/message"): case r.Method == http.MethodGet && strings.HasSuffix(r.URL.Path, "/message"):
// поллинг прогресса: голый массив [{info, parts}]. // v2 поллинг: {data:[Session.Message]} (новейшие первыми).
id := strings.TrimSuffix(strings.TrimPrefix(r.URL.Path, "/session/"), "/message") id := sessionID(r.URL.Path, "/message")
mu.Lock() mu.Lock()
agent := sessions[id] agent := sessions[id]
mu.Unlock() mu.Unlock()
writeJSON(w, []map[string]any{{"info": map[string]any{"role": "assistant"}, "parts": []map[string]any{{"type": "text", "text": verdictFor(agent)}}}}) writeJSON(w, map[string]any{"data": []map[string]any{assistantMsg(id, agent)}})
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/abort"): case r.Method == http.MethodGet && r.URL.Path == "/api/session/active":
writeJSON(w, map[string]any{"ok": true}) // сессий в активных дренажах нет → ответ завершён.
writeJSON(w, map[string]any{"data": map[string]any{}})
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/interrupt"):
writeJSON(w, map[string]any{"data": map[string]any{"ok": true}})
default: default:
http.NotFound(w, r) http.NotFound(w, r)
@@ -143,6 +176,7 @@ func e2eAssemble(t *testing.T) (*App, string, *e2eChannel) {
CoreCtx: coreCtx, CoreCtx: coreCtx,
} }
a.Router = chat.NewRouter(a.handleIncoming) a.Router = chat.NewRouter(a.handleIncoming)
fake.router = a.Router
if err := a.Router.Attach(fake); err != nil { if err := a.Router.Attach(fake); err != nil {
t.Fatalf("Attach fake channel: %v", err) t.Fatalf("Attach fake channel: %v", err)
} }
@@ -168,6 +202,7 @@ func e2eAssemble(t *testing.T) (*App, string, *e2eChannel) {
type e2eChannel struct { type e2eChannel struct {
onMsg chat.Handler onMsg chat.Handler
sent []chat.Message sent []chat.Message
router *chat.Router
} }
func (c *e2eChannel) Run(_ context.Context) error { return nil } func (c *e2eChannel) Run(_ context.Context) error { return nil }
@@ -184,10 +219,21 @@ func (c *e2eChannel) Ask(_ context.Context, _ chat.Address, m chat.Message) erro
func (c *e2eChannel) Close() error { return nil } func (c *e2eChannel) Close() error { return nil }
// deliver отправляет входящее сообщение через роутер: ставит маршрут // deliver отправляет входящее сообщение через роутер: ставит маршрут
// пользователя и вызывает app.handleIncoming (как в проде). // пользователя и вызывает app.handleIncoming (как в проде). Так как роутер
// обрабатывает входящие асинхронно (процессор-горутина), deliver ждёт, пока
// обработка события завершится, — иначе тесты (сразу читающие состояние БД)
// гоняются с обработчиком.
func (c *e2eChannel) deliver(uid chat.UserID, text string) { func (c *e2eChannel) deliver(uid chat.UserID, text string) {
if c.onMsg != nil { if c.onMsg == nil {
return
}
target := int64(0)
if c.router != nil {
target = c.router.Processed() + 1
}
c.onMsg(chat.Incoming{UserID: uid, Address: chat.Address("u://" + string(uid)), Msg: chat.Message{Text: text}, Channel: c}) c.onMsg(chat.Incoming{UserID: uid, Address: chat.Address("u://" + string(uid)), Msg: chat.Message{Text: text}, Channel: c})
if c.router != nil && !c.router.WaitProcessed(target) {
panic("e2e: роутер не обработал входящее за 30s")
} }
} }

15
internal/app/ui_cgo.go Normal file
View File

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

11
internal/app/ui_noui.go Normal file
View File

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

View File

@@ -4,6 +4,8 @@ import (
"context" "context"
"fmt" "fmt"
"sync" "sync"
"sync/atomic"
"time"
) )
// Router — единый диспетчер входящих из всех каналов и маршрутизатор исходящих. // Router — единый диспетчер входящих из всех каналов и маршрутизатор исходящих.
@@ -24,6 +26,10 @@ type Router struct {
// long-poll цикл канала (Telegram) не блокируется на время долгого // long-poll цикл канала (Telegram) не блокируется на время долгого
// вызова аналитика и продолжает принимать новые сообщения. // вызова аналитика и продолжает принимать новые сообщения.
incoming chan Incoming incoming chan Incoming
// processed — число обработанных воркером событий (для синхронизации
// тестов с асинхронной очередью: WaitProcessed ждёт обработку события).
processed atomic.Int64
} }
// NewRouter создаёт роутер. onUserMsg — колбэк обработки входящего. // NewRouter создаёт роутер. onUserMsg — колбэк обработки входящего.
@@ -46,9 +52,26 @@ func NewRouter(onUserMsg func(Incoming)) *Router {
func (r *Router) processLoop() { func (r *Router) processLoop() {
for inc := range r.incoming { for inc := range r.incoming {
r.onUserMsg(inc) r.onUserMsg(inc)
r.processed.Add(1)
} }
} }
// Processed возвращает число обработанных воркером входящих событий.
func (r *Router) Processed() int64 { return r.processed.Load() }
// WaitProcessed ждёт, пока воркер обработает не меньше target событий
// (для синхронизации с асинхронной очередью в тестах).
func (r *Router) WaitProcessed(target int64) bool {
deadline := time.Now().Add(30 * time.Second)
for r.processed.Load() < target {
if time.Now().After(deadline) {
return false
}
time.Sleep(2 * time.Millisecond)
}
return true
}
// Attach регистрирует канал и подключает его к обработчику входящих. // Attach регистрирует канал и подключает его к обработчику входящих.
// Возвращает ошибку только при пустом канале (nil). // Возвращает ошибку только при пустом канале (nil).
func (r *Router) Attach(ch Channel) error { func (r *Router) Attach(ch Channel) error {

View File

@@ -251,14 +251,17 @@ func TestResolveExePaths_AbsoluteKept(t *testing.T) {
t.Setenv("TG_TOKEN", "tok") t.Setenv("TG_TOKEN", "tok")
t.Setenv("TG_CHAT_ID", "42") t.Setenv("TG_CHAT_ID", "42")
absDB := filepath.Join(string(filepath.Separator), "data", "ratatoskr.db") // абсолютный для текущей ОС // абсолютные пути «для текущей ОС»: на Windows слэш-относительный путь
absWt := filepath.Join(string(filepath.Separator), "worktrees") // (\data\...) НЕ является абсолютным — нужен корень тома (C:\data\...).
root := filepath.VolumeName(os.TempDir()) + string(filepath.Separator)
absDB := filepath.Join(root, "data", "ratatoskr.db")
absWt := filepath.Join(root, "worktrees")
yaml := `telegram: yaml := `telegram:
username: "${TG_TOKEN}" username: "${TG_TOKEN}"
chat_id: "${TG_CHAT_ID}" chat_id: "${TG_CHAT_ID}"
paths: paths:
db: "` + absDB + `" db: '` + absDB + `'
worktree: "` + absWt + `" worktree: '` + absWt + `'
` `
cfg, err := Load(writeCfg(t, yaml)) cfg, err := Load(writeCfg(t, yaml))
if err != nil { if err != nil {

View File

@@ -5,6 +5,7 @@ import (
"fmt" "fmt"
"strings" "strings"
"github.com/kamelion/ratatoskr-go/internal/events"
"github.com/kamelion/ratatoskr-go/internal/storage" "github.com/kamelion/ratatoskr-go/internal/storage"
) )
@@ -18,6 +19,27 @@ type Core struct {
MaxQuestionsPerTurn int MaxQuestionsPerTurn int
// Live — опциональный просмотр живой сессии задачи (для /status N). // Live — опциональный просмотр живой сессии задачи (для /status N).
Live LiveProber 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). // LiveProber — абстракция за журналом живых сессий (реализация — *opencode.LiveRegistry).
@@ -117,11 +139,10 @@ func (c *Core) handleStart(ctx context.Context, taskID int64) (Result, error) {
if err != nil { if err != nil {
return Result{}, err return Result{}, err
} }
task.Status = storage.StatusCollecting if err := c.setStatus(ctx, task, storage.StatusCollecting); err != nil {
if err := c.Store.ClearHistory(ctx, task.ID); err != nil {
return Result{}, err 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{}, err
} }
return Result{ return Result{
@@ -137,8 +158,7 @@ func (c *Core) handleCancel(ctx context.Context, taskID int64) (Result, error) {
if err != nil { if err != nil {
return Result{}, err return Result{}, err
} }
task.Status = storage.StatusCancelled if err := c.setStatus(ctx, task, storage.StatusCancelled); err != nil {
if err := c.Store.UpdateTask(ctx, task); err != nil {
return Result{}, err return Result{}, err
} }
return Result{ return Result{
@@ -202,11 +222,10 @@ func (c *Core) handleRetry(ctx context.Context, rest string) (Result, error) {
Status: task.Status, Status: task.Status,
}, nil }, nil
} }
task.Status = storage.StatusCollecting if err := c.setStatus(ctx, task, storage.StatusCollecting); err != nil {
if err := c.Store.ClearHistory(ctx, id); err != nil {
return Result{}, err 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{}, err
} }
return Result{ return Result{
@@ -272,8 +291,7 @@ func (c *Core) notFoundReply(ctx context.Context, id int64, err error) (Result,
// handleConsent одобряет задачу: ready → approved (финальное одобрение, // handleConsent одобряет задачу: ready → approved (финальное одобрение,
// после которого воркер забирает задачу на выполнение). // после которого воркер забирает задачу на выполнение).
func (c *Core) handleConsent(ctx context.Context, task *storage.Task) (Result, error) { func (c *Core) handleConsent(ctx context.Context, task *storage.Task) (Result, error) {
task.Status = storage.StatusApproved if err := c.setStatus(ctx, task, storage.StatusApproved); err != nil {
if err := c.Store.UpdateTask(ctx, task); err != nil {
return Result{}, err return Result{}, err
} }
return Result{ return Result{
@@ -285,6 +303,7 @@ func (c *Core) handleConsent(ctx context.Context, task *storage.Task) (Result, e
// handleEdit — правка черновика в фазе ready → снова сбор + аналитик. // handleEdit — правка черновика в фазе ready → снова сбор + аналитик.
func (c *Core) handleEdit(ctx context.Context, task *storage.Task, text string) (Result, error) { func (c *Core) handleEdit(ctx context.Context, task *storage.Task, text string) (Result, error) {
from := task.Status
task.Status = storage.StatusCollecting task.Status = storage.StatusCollecting
if err := c.Store.AppendHistory(ctx, task.ID, "user", text); err != nil { if err := c.Store.AppendHistory(ctx, task.ID, "user", text); err != nil {
return Result{}, err 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 { if err := c.Store.UpdateTask(ctx, task); err != nil {
return Result{}, err return Result{}, err
} }
c.publish(events.TaskStatusChanged{ID: task.ID, From: from, To: task.Status})
return c.runDecide(ctx, task, false) return c.runDecide(ctx, task, false)
} }
@@ -300,10 +320,12 @@ func (c *Core) handleTurn(ctx context.Context, task *storage.Task, text string)
// переход draft→collecting нужно персистить до вызова аналитика, // переход draft→collecting нужно персистить до вызова аналитика,
// иначе propose сделает draft→ready (невалидно) // иначе propose сделает draft→ready (невалидно)
if task.Status != storage.StatusCollecting { if task.Status != storage.StatusCollecting {
from := task.Status
task.Status = storage.StatusCollecting task.Status = storage.StatusCollecting
if err := c.Store.UpdateTask(ctx, task); err != nil { if err := c.Store.UpdateTask(ctx, task); err != nil {
return Result{}, err return Result{}, err
} }
c.publish(events.TaskStatusChanged{ID: task.ID, From: from, To: task.Status})
} }
if text != "" { if text != "" {
if err := c.Store.AppendHistory(ctx, task.ID, "user", text); err != nil { 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 { switch decision.Phase {
case "abort": case "abort":
task.Status = storage.StatusAborted if err := c.setStatus(ctx, task, storage.StatusAborted); err != nil {
if err := c.Store.UpdateTask(ctx, task); err != nil {
return Result{}, err return Result{}, err
} }
reply := decision.ChatReply reply := decision.ChatReply
@@ -363,15 +384,13 @@ func (c *Core) runDecide(ctx context.Context, task *storage.Task, force bool) (R
applyDraft(task, decision.Draft) applyDraft(task, decision.Draft)
// E1: propose/ready без репозиториев → остаёмся в сборе, просим уточнить. // E1: propose/ready без репозиториев → остаёмся в сборе, просим уточнить.
if len(task.EffectiveRepos()) == 0 { if len(task.EffectiveRepos()) == 0 {
task.Status = storage.StatusCollecting if err := c.setStatus(ctx, task, storage.StatusCollecting); err != nil {
if err := c.Store.UpdateTask(ctx, task); err != nil {
return Result{}, err return Result{}, err
} }
reply := decChatReply(decision, "Укажи, в каком репозитории(ях) вести работу.") reply := decChatReply(decision, "Укажи, в каком репозитории(ях) вести работу.")
return Result{Reply: reply, TaskID: task.ID, Status: task.Status}, nil return Result{Reply: reply, TaskID: task.ID, Status: task.Status}, nil
} }
task.Status = storage.StatusReady if err := c.setStatus(ctx, task, storage.StatusReady); err != nil {
if err := c.Store.UpdateTask(ctx, task); err != nil {
return Result{}, err return Result{}, err
} }
return Result{ return Result{
@@ -382,8 +401,7 @@ func (c *Core) runDecide(ctx context.Context, task *storage.Task, force bool) (R
default: // "ask" default: // "ask"
applyDraft(task, decision.Draft) applyDraft(task, decision.Draft)
task.Status = storage.StatusCollecting if err := c.setStatus(ctx, task, storage.StatusCollecting); err != nil {
if err := c.Store.UpdateTask(ctx, task); err != nil {
return Result{}, err return Result{}, err
} }
reply := buildAskReply(decision, c.MaxQuestionsPerTurn) reply := buildAskReply(decision, c.MaxQuestionsPerTurn)

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

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

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

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

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

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

96
internal/events/hub.go Normal file
View File

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

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

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

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

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

View File

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

View File

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

77
internal/model/status.go Normal file
View File

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

View File

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

View File

@@ -11,29 +11,28 @@ import (
"time" "time"
) )
// Client — HTTP-взаимодействие с одним opencode serve (режим API). // Client — HTTP-взаимодействие с одним opencode serve (v2 HTTP API).
// //
// Ходит по experimental HTTP API opencode serve (пути БЕЗ префикса /api): // Пути v2 начинаются с префикса /api (см. README, минимальная версия opencode):
// - POST /session создать сессию → голая Session {id} // - POST /api/session создать сессию {model:{...}} → {data: Session.Info}
// - POST /session/{id}/message отправить промпт {parts:[{type:"text"}]} → // - POST /api/session/{id}/prompt отправить промпт {prompt:{text}} →
// блокирует и возвращает {info,parts}; вердикт из parts // НЕБЛОКИРУЮЩЕ (admit) → {data: Admitted}
// - GET /session/{id}/message история → голый массив [{info, parts}] (для прогресса) // - GET /api/session/{id}/message?order=desc → {data:[Message,...]}
// - POST /session/{id}/abort прервать выполняющийся ответ // - POST /api/session/{id}/interrupt прервать активный ответ (204)
// - GET /api/session/active активные дренажи → {data:{sessionID:...}}
// //
// Вердикт собирается из parts[] ответа на POST /message: текст тех частей, // Prompt не блокирует: вердикт собирается поллингом из content[].type=="text"
// где type == "text". // новых assistant-сообщений (см. Runner.awaitVerdict).
type Client struct { type Client struct {
BaseURL string // http://host:port (без завершающего слеша) BaseURL string // http://host:port (без завершающего слеша)
Password string // basic auth (username "opencode") Password string // basic auth (username "opencode")
Debug bool // включать отладочные логи API-вызовов (log.level=debug) Debug bool // включать отладочные логи API-вызовов (log.level=debug)
http *http.Client // для быстрых операций (create/messages/abort) http *http.Client // единый клиент: все операции быстрые (нет блокирующего Send)
httpSend *http.Client // для блокирующего Send — без жёсткого таймаута,
// отменяется только через контекст (idle/hard)
} }
// ClientErr — классы ошибок клиента. // ClientErr — классы ошибок клиента.
type ClientErr struct { type ClientErr struct {
Op string // "connect" | "create" | "prompt" | "messages" | "abort" Op string // "connect" | "create" | "prompt" | "messages" | "active" | "abort"
Err error Err error
} }
@@ -44,19 +43,10 @@ func (c *Client) defaults() {
if c.http == nil { if c.http == nil {
c.http = &http.Client{Timeout: 30 * time.Second} c.http = &http.Client{Timeout: 30 * time.Second}
} }
if c.httpSend == nil {
c.httpSend = &http.Client{}
}
} }
// do выполняет запрос через c.http (с таймаутом 30s) и возвращает тело при 2xx. // do выполняет запрос через c.http и возвращает тело при 2xx.
func (c *Client) do(ctx context.Context, method, path, op string, body []byte) ([]byte, error) { func (c *Client) do(ctx context.Context, method, path, op string, body []byte) ([]byte, error) {
c.defaults()
return c.doHTTP(ctx, method, path, op, body, c.http)
}
// doHTTP — общая реализация запроса; hc — клиент, которым выполняется запрос.
func (c *Client) doHTTP(ctx context.Context, method, path, op string, body []byte, hc *http.Client) ([]byte, error) {
c.defaults() c.defaults()
var rd io.Reader var rd io.Reader
if body != nil { if body != nil {
@@ -78,7 +68,7 @@ func (c *Client) doHTTP(ctx context.Context, method, path, op string, body []byt
log.Printf("opencode api %s request body: %s", op, truncateStr(string(body), 5000)) log.Printf("opencode api %s request body: %s", op, truncateStr(string(body), 5000))
} }
} }
resp, err := hc.Do(req) resp, err := c.http.Do(req)
if err != nil { if err != nil {
return nil, &ClientErr{Op: "connect", Err: err} return nil, &ClientErr{Op: "connect", Err: err}
} }
@@ -99,116 +89,238 @@ func (c *Client) doHTTP(ctx context.Context, method, path, op string, body []byt
return b, nil return b, nil
} }
// CreateSession создаёт новую сессию и возвращает её id. // ModelRef — ссылка на модель (аналог v2 Model.Ref: {providerID, id, variant?}).
func (c *Client) CreateSession(ctx context.Context, title string) (string, error) { // providerID — имя провайдера из конфига opencode, id — идентификатор модели.
body := map[string]string{} type ModelRef struct {
if title != "" { ProviderID string `json:"providerID"`
body["title"] = title ID string `json:"id"`
Variant string `json:"variant,omitempty"`
} }
b, _ := json.Marshal(body)
raw, err := c.do(ctx, http.MethodPost, "/session", "create", b) // String возвращает каноничное представление "provider/id[/variant]".
func (m *ModelRef) String() string {
if m == nil {
return ""
}
if m.Variant != "" {
return m.ProviderID + "/" + m.ID + "/" + m.Variant
}
return m.ProviderID + "/" + m.ID
}
// CreateSession создаёт новую сессию и возвращает её id. model != nil —
// хардпин модели (top-level "model" из конфига opencode), чтобы не зависеть
// от fallback-логики выбора модели в самом opencode.
func (c *Client) CreateSession(ctx context.Context, model *ModelRef) (string, error) {
payload := map[string]any{}
if model != nil {
payload["model"] = model
}
body, _ := json.Marshal(payload)
raw, err := c.do(ctx, http.MethodPost, "/api/session", "create", body)
if err != nil { if err != nil {
return "", err return "", err
} }
// experimental: ответ — голая Session (без обёртки {data}).
var out struct { var out struct {
Data struct {
ID string `json:"id"` ID string `json:"id"`
} `json:"data"`
} }
if err := json.Unmarshal(raw, &out); err != nil { if err := json.Unmarshal(raw, &out); err != nil {
return "", &ClientErr{Op: "create", Err: fmt.Errorf("невалидный ответ: %v", err)} return "", &ClientErr{Op: "create", Err: fmt.Errorf("невалидный ответ: %v", err)}
} }
if out.ID == "" { if out.Data.ID == "" {
return "", &ClientErr{Op: "create", Err: fmt.Errorf("пустой id сессии")} return "", &ClientErr{Op: "create", Err: fmt.Errorf("пустой id сессии")}
} }
return out.ID, nil return out.Data.ID, nil
} }
// Send отправляет промпт в сессию, БЛОКИРУЯСЬ до завершения ответа, и // Admitted — результат admit промпта (SessionInput.Admitted).
// возвращает вердикт (текст text-частей из parts). Отмена — только через ctx type Admitted struct {
// (используется отдельный клиент без жёсткого таймаута; idle/hard в Runner'е ID string // id user-сообщения
// отменяют контекст, что прерывает этот запрос). TimeCreated int64 // epoch ms создания промпта (граница «новых» ответов)
func (c *Client) Send(ctx context.Context, sessionID, prompt string) (string, error) {
payload := map[string]any{
"parts": []map[string]string{{"type": "text", "text": prompt}},
} }
b, _ := json.Marshal(payload)
c.defaults() // Prompt неблокирующе отправляет промпт в сессию (durable admit) и возвращает
raw, err := c.doHTTP(ctx, http.MethodPost, "/session/"+sessionID+"/message", "prompt", b, c.httpSend) // границу времени, с которой следует считать assistant-сообщения «новыми».
func (c *Client) Prompt(ctx context.Context, sessionID, prompt string) (*Admitted, error) {
payload := map[string]any{
"prompt": map[string]string{"text": prompt},
}
body, _ := json.Marshal(payload)
raw, err := c.do(ctx, http.MethodPost, "/api/session/"+sessionID+"/prompt", "prompt", body)
if err != nil { if err != nil {
return "", err return nil, err
} }
var out struct { var out struct {
Parts []part `json:"parts"` Data struct {
ID string `json:"id"`
TimeCreated int64 `json:"timeCreated"`
} `json:"data"`
} }
if err := json.Unmarshal(raw, &out); err != nil { if err := json.Unmarshal(raw, &out); err != nil {
return "", &ClientErr{Op: "prompt", Err: fmt.Errorf("невалидный ответ: %v", err)} return nil, &ClientErr{Op: "prompt", Err: fmt.Errorf("невалидный ответ: %v", err)}
} }
var buf bytes.Buffer if out.Data.ID == "" {
for _, p := range out.Parts { return nil, &ClientErr{Op: "prompt", Err: fmt.Errorf("пустой id промпта в ответе")}
if p.Type == "text" && p.Text != "" {
if buf.Len() > 0 {
buf.WriteString("\n")
} }
buf.WriteString(p.Text) return &Admitted{ID: out.Data.ID, TimeCreated: out.Data.TimeCreated}, nil
}
}
if buf.Len() == 0 {
return "", &ClientErr{Op: "prompt", Err: fmt.Errorf("нет text-части в ответе")}
}
return stripFence(buf.String()), nil
} }
// Abort прерывает выполняющийся ответ сессии. // v2Message — минимальная проекция Session.Message (tagged union: тип в "type").
func (c *Client) Abort(ctx context.Context, sessionID string) error { // Поле "role" в v2 отсутствует; assistant определяется по type=="assistant".
_, err := c.do(ctx, http.MethodPost, "/session/"+sessionID+"/abort", "abort", nil) type v2Message struct {
return err ID string `json:"id"`
Type string `json:"type"` // "assistant" | "user" | "tool" | "system" | ...
Content []v2Part `json:"content"`
Model *ModelRef `json:"model"`
Finish string `json:"finish,omitempty"`
Error *v2Error `json:"error,omitempty"`
Time v2Time `json:"time"`
} }
// part — минимальная часть сообщения (из parts[]). type v2Part struct {
type part struct {
Type string `json:"type"` // "text" | "reasoning" | "tool" | ... Type string `json:"type"` // "text" | "reasoning" | "tool" | ...
Text string `json:"text"` Text string `json:"text"`
} }
// message — элемент голого массива из GET /session/{id}/message. type v2Time struct {
type message struct { Created *int64 `json:"created"`
Info struct { Completed *int64 `json:"completed"`
Role string `json:"role"` // "assistant" | "user" | ...
} `json:"info"`
Parts []part `json:"parts"`
} }
// messages возвращает сырые сообщения сессии (для поллинга прогресса). type v2Error struct {
func (c *Client) messages(ctx context.Context, sessionID string) ([]message, error) { Type string `json:"type"`
raw, err := c.do(ctx, http.MethodGet, "/session/"+sessionID+"/message", "messages", nil) Message string `json:"message"`
}
// finished — завершено ли assistant-сообщение (ответ агента закончен).
func (m *v2Message) finished() bool {
if m == nil {
return false
}
if m.Error != nil {
return true
}
if m.Finish != "" {
return true
}
return m.Time.Completed != nil && *m.Time.Completed > 0
}
// Messages возвращает сообщения сессии (новейшие первыми, до 200 за запрос).
func (c *Client) Messages(ctx context.Context, sessionID string) ([]v2Message, error) {
raw, err := c.do(ctx, http.MethodGet, "/api/session/"+sessionID+"/message?order=desc&limit=200", "messages", nil)
if err != nil { if err != nil {
return nil, err return nil, err
} }
var out []message var out struct {
Data []v2Message `json:"data"`
}
if err := json.Unmarshal(raw, &out); err != nil { if err := json.Unmarshal(raw, &out); err != nil {
return nil, err return nil, &ClientErr{Op: "messages", Err: fmt.Errorf("невалидный ответ: %v", err)}
} }
return out, nil return out.Data, nil
} }
// textCount считает число text-частей в assistant-сообщениях (для progress). // Active возвращает true, если сессия ещё обрабатывается (есть в активных
func (c *Client) textCount(ctx context.Context, sessionID string) (int, error) { // дренажах этого serve). Сессии вне списка считаются завершёнными.
msgs, err := c.messages(ctx, sessionID) func (c *Client) Active(ctx context.Context, sessionID string) (bool, error) {
raw, err := c.do(ctx, http.MethodGet, "/api/session/active", "active", nil)
if err != nil { if err != nil {
return 0, err return false, err
} }
n := 0 var out struct {
for _, m := range msgs { Data map[string]json.RawMessage `json:"data"`
if m.Info.Role != "assistant" { }
if err := json.Unmarshal(raw, &out); err != nil {
return false, &ClientErr{Op: "active", Err: fmt.Errorf("невалидный ответ: %v", err)}
}
if out.Data == nil {
return false, nil
}
_, ok := out.Data[sessionID]
return ok, nil
}
// Interrupt прерывает активный ответ сессии (аналог v1 abort).
func (c *Client) Interrupt(ctx context.Context, sessionID string) error {
_, err := c.do(ctx, http.MethodPost, "/api/session/"+sessionID+"/interrupt", "abort", nil)
return err
}
// assistantSince фильтрует assistant-сообщения, созданные не раньше since
// (порядок сохраняется — как пришёл из API, новейшие первыми).
func assistantSince(msgs []v2Message, since int64) []*v2Message {
out := make([]*v2Message, 0, len(msgs))
for i := range msgs {
m := &msgs[i]
if m.Type != "assistant" {
continue continue
} }
for _, p := range m.Parts { if m.Time.Created == nil || *m.Time.Created < since {
continue
}
out = append(out, m)
}
return out
}
// textParts считает text-парты в одном assistant-сообщении (для прогресса).
func textParts(m *v2Message) int {
n := 0
for _, p := range m.Content {
if p.Type == "text" && p.Text != "" { if p.Type == "text" && p.Text != "" {
n++ n++
} }
} }
return n
} }
return n, nil
// newestAssistant возвращает самое новое assistant-сообщение (из фильтра) и
// суммарное число text-партов. since — граница времени (epoch ms).
func newestAssistant(msgs []v2Message, since int64) (*v2Message, int) {
ass := assistantSince(msgs, since)
var newest *v2Message
count := 0
for _, m := range ass {
count += textParts(m)
if newest == nil || *m.Time.Created > *newest.Time.Created {
newest = m
}
}
return newest, count
}
// progressOf — «живой» прогресс новых assistant-сообщений: число контент-партов
// (text/reasoning/tool) + суммарная длина их текста. Растёт во время стриминга,
// когда один и тот же парт увеличивается (и при reasoning), — это и есть
// сигнал, что LLM работает, а не висит.
func progressOf(msgs []v2Message, since int64) (parts, textLen int) {
for _, m := range assistantSince(msgs, since) {
for _, p := range m.Content {
parts++
if p.Type == "text" || p.Type == "reasoning" {
textLen += len(p.Text)
}
}
}
return
}
// assistantText объединяет text-парты новых assistant-сообщений в хронологическом
// порядке (сообщения приходят новейшими первыми → идём с конца).
func assistantText(msgs []v2Message, since int64) []string {
ass := assistantSince(msgs, since)
texts := make([]string, 0, len(ass))
for i := len(ass) - 1; i >= 0; i-- {
for _, p := range ass[i].Content {
if p.Type == "text" && p.Text != "" {
texts = append(texts, p.Text)
}
}
}
return texts
} }
func truncateStr(s string, n int) string { func truncateStr(s string, n int) string {

View File

@@ -6,22 +6,45 @@ import (
"errors" "errors"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"strings"
"testing" "testing"
"time" "time"
) )
// fakeAPIServer — минимальный фейк opencode serve experimental HTTP API // fakeAPIServer — минимальный фейк opencode serve v2 HTTP API (пути /api/*).
// (пути БЕЗ префикса /api). //
// Сценарии:
// - нормальный: Prompt ставит active=false и в messages кладётся финальное
// assistant-сообщение (verdictText) → Runner собирает вердикт;
// - blockPrompt: «агент завис» — active=true всегда, сообщений нет → idle abort;
// - failCreate / failMessages — имитация ошибок;
// - growStream: стрим одного растущего парта — текст/reasoning растёт с
// каждым опросом GET /message (streamPolls раз), active=true, затем
// active=false + финальное завершённое сообщение.
type fakeAPIServer struct { type fakeAPIServer struct {
messages []message sessionID string
created bool
active bool
blockPrompt bool
messages []v2Message
verdictText string
failCreate bool failCreate bool
verdictParts []part // ответ на POST /session/{id}/message (вердикт) failMessages bool
blockPrompt bool // POST /message блокируется до отмены ctx (эмуляция зависания) createdModel *ModelRef // модель, полученная на POST /api/session
promptCalls int
// streamGrow: стрим одного растущего парта — текст/reasoning растёт с
// каждым опросом GET /message, active=true, пока messageCalls не дойдёт до
// streamPolls; затем active=false + финальное завершённое сообщение.
streamGrow bool
streamReasoning bool // растущий парт — reasoning вместо text
streamPolls int // сколько опросов длится «стрим» до завершения
messageCalls int
} }
func (f *fakeAPIServer) handler() http.Handler { func (f *fakeAPIServer) handler() http.Handler {
mux := http.NewServeMux() mux := http.NewServeMux()
mux.HandleFunc("/session", func(w http.ResponseWriter, r *http.Request) { mux.HandleFunc("/api/session", func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost { if r.Method != http.MethodPost {
w.WriteHeader(http.StatusMethodNotAllowed) w.WriteHeader(http.StatusMethodNotAllowed)
return return
@@ -30,50 +53,106 @@ func (f *fakeAPIServer) handler() http.Handler {
http.Error(w, "boom", http.StatusInternalServerError) http.Error(w, "boom", http.StatusInternalServerError)
return return
} }
// experimental: голая Session (без обёртки {data}). var in struct {
writeJSON(w, map[string]any{"id": "sess-fake"}) Model *ModelRef `json:"model"`
}
_ = json.NewDecoder(r.Body).Decode(&in)
f.createdModel = in.Model
f.sessionID = "sess-fake"
f.created = true
writeJSON(w, map[string]any{"data": map[string]any{"id": "sess-fake"}})
}) })
mux.HandleFunc("/session/{id}/abort", func(w http.ResponseWriter, r *http.Request) { mux.HandleFunc("/api/session/active", func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
w.WriteHeader(http.StatusMethodNotAllowed)
return
}
data := map[string]any{}
if f.active && f.sessionID != "" {
data[f.sessionID] = map[string]any{"type": "running"}
}
writeJSON(w, map[string]any{"data": data})
})
mux.HandleFunc("/api/session/{id}/prompt", func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost { if r.Method != http.MethodPost {
w.WriteHeader(http.StatusMethodNotAllowed) w.WriteHeader(http.StatusMethodNotAllowed)
return return
} }
writeJSON(w, map[string]any{}) f.promptCalls++
})
mux.HandleFunc("/session/{id}/message", func(w http.ResponseWriter, r *http.Request) {
switch r.Method {
case http.MethodPost:
if f.blockPrompt { if f.blockPrompt {
// Эмуляция «зависшего» агента: ответ приходит позже idle-таймаута, // «зависший» агент: активен, но сообщений не появляется.
// но handler всё равно завершится, чтобы не блокировать shutdown. f.active = true
select { } else {
case <-r.Context().Done(): f.active = false
case <-time.After(2 * time.Second):
} }
w.WriteHeader(http.StatusRequestTimeout) writeJSON(w, map[string]any{"data": map[string]any{
return "id": "msg_1",
} "sessionID": f.sessionID,
// блокирующий ответ: {info, parts}, где вердикт — text-части. "timeCreated": time.Now().UnixMilli(),
info := map[string]any{"role": "assistant"} }})
parts := f.verdictParts })
if parts == nil { mux.HandleFunc("/api/session/{id}/interrupt", func(w http.ResponseWriter, r *http.Request) {
parts = []part{} if r.Method != http.MethodPost {
}
writeJSON(w, map[string]any{"info": info, "parts": parts})
case http.MethodGet:
// голый массив [{info, parts}].
if f.messages == nil {
writeJSON(w, []message{})
return
}
writeJSON(w, f.messages)
default:
w.WriteHeader(http.StatusMethodNotAllowed) w.WriteHeader(http.StatusMethodNotAllowed)
return
} }
f.active = false
w.WriteHeader(http.StatusNoContent)
})
mux.HandleFunc("/api/session/{id}/message", func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
w.WriteHeader(http.StatusMethodNotAllowed)
return
}
if f.failMessages {
http.Error(w, "db error", http.StatusInternalServerError)
return
}
if f.streamGrow {
f.messageCalls++
now := time.Now().UnixMilli()
done := f.messageCalls >= f.streamPolls
msg := v2Message{ID: "msg_stream", Type: "assistant", Time: v2Time{Created: &now}}
switch {
case done:
msg.Content = []v2Part{{Type: "text", Text: "done-stream"}}
msg.Finish = "end_turn"
msg.Time.Completed = &now
f.active = false
case f.streamReasoning:
msg.Content = []v2Part{{Type: "reasoning", Text: strings.Repeat("r", f.messageCalls)}}
f.active = true
default:
msg.Content = []v2Part{{Type: "text", Text: strings.Repeat("x", f.messageCalls)}}
f.active = true
}
writeJSON(w, map[string]any{"data": []v2Message{msg}})
return
}
msgs := f.messages
if msgs == nil && f.verdictText != "" && !f.blockPrompt {
msgs = []v2Message{f.assistantMsg(f.verdictText)}
}
if msgs == nil {
msgs = []v2Message{}
}
writeJSON(w, map[string]any{"data": msgs})
}) })
return mux return mux
} }
// assistantMsg строит завершённое assistant-сообщение с text-партом.
func (f *fakeAPIServer) assistantMsg(text string) v2Message {
now := time.Now().UnixMilli()
return v2Message{
ID: "msg_a",
Type: "assistant",
Content: []v2Part{{Type: "text", Text: text}},
Finish: "end_turn",
Time: v2Time{Created: &now, Completed: &now},
}
}
func writeJSON(w http.ResponseWriter, v any) { func writeJSON(w http.ResponseWriter, v any) {
w.Header().Set("Content-Type", "application/json") w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(v) _ = json.NewEncoder(w).Encode(v)
@@ -88,61 +167,135 @@ func fakeClient(t *testing.T, f *fakeAPIServer) *Client {
} }
func TestClient_CreateSession(t *testing.T) { func TestClient_CreateSession(t *testing.T) {
c := fakeClient(t, &fakeAPIServer{}) f := &fakeAPIServer{}
id, err := c.CreateSession(context.Background(), "ratatoskr-analyst") c := fakeClient(t, f)
id, err := c.CreateSession(context.Background(), nil)
if err != nil { if err != nil {
t.Fatalf("CreateSession err: %v", err) t.Fatalf("CreateSession err: %v", err)
} }
if id != "sess-fake" { if id != "sess-fake" {
t.Errorf("id = %q, want sess-fake", id) t.Errorf("id = %q, want sess-fake", id)
} }
if f.createdModel != nil {
t.Errorf("createdModel = %+v, want nil", f.createdModel)
}
}
func TestClient_CreateSessionHardpinsModel(t *testing.T) {
want := &ModelRef{ProviderID: "tokentool", ID: "deepseek/deepseek-v4-flash-0731"}
f := &fakeAPIServer{}
c := fakeClient(t, f)
if _, err := c.CreateSession(context.Background(), want); err != nil {
t.Fatalf("CreateSession err: %v", err)
}
if f.createdModel == nil || f.createdModel.ProviderID != want.ProviderID || f.createdModel.ID != want.ID {
t.Errorf("createdModel = %+v, want %+v", f.createdModel, want)
}
} }
func TestClient_CreateSessionFail(t *testing.T) { func TestClient_CreateSessionFail(t *testing.T) {
c := fakeClient(t, &fakeAPIServer{failCreate: true}) c := fakeClient(t, &fakeAPIServer{failCreate: true})
if _, err := c.CreateSession(context.Background(), "x"); err == nil { if _, err := c.CreateSession(context.Background(), nil); err == nil {
t.Fatal("CreateSession должен упасть при 500, а не nil") t.Fatal("CreateSession должен упасть при 500, а не nil")
} }
} }
func TestClient_Send(t *testing.T) { func TestClient_Prompt(t *testing.T) {
c := fakeClient(t, &fakeAPIServer{ c := fakeClient(t, &fakeAPIServer{})
verdictParts: []part{{Type: "text", Text: `{"phase":"ready"}`}}, adm, err := c.Prompt(context.Background(), "sess-fake", "почини x")
})
vd, err := c.Send(context.Background(), "sess-fake", "почини x")
if err != nil { if err != nil {
t.Fatalf("Send err: %v", err) t.Fatalf("Prompt err: %v", err)
} }
if vd != `{"phase":"ready"}` { if adm.ID != "msg_1" {
t.Errorf("verdict = %q, want вердикт модели", vd) t.Errorf("adm.ID = %q, want msg_1", adm.ID)
}
if adm.TimeCreated == 0 {
t.Error("adm.TimeCreated = 0, want epoch ms")
} }
} }
func TestClient_SendNoText(t *testing.T) { func TestClient_Messages(t *testing.T) {
c := fakeClient(t, &fakeAPIServer{}) // нет text-части в ответе now := time.Now().UnixMilli()
if _, err := c.Send(context.Background(), "sess-fake", "почини x"); err == nil { f := &fakeAPIServer{messages: []v2Message{{
t.Fatal("Send должен упасть, когда нет text-части") ID: "msg_a", Type: "assistant",
} else { Content: []v2Part{{Type: "text", Text: "a"}, {Type: "reasoning", Text: "x"}},
var ce *ClientErr Finish: "end_turn",
if !errors.As(err, &ce) { Time: v2Time{Created: &now, Completed: &now},
t.Errorf("ожидался *ClientErr, got %T", err)
}
}
}
func TestClient_textCount(t *testing.T) {
f := &fakeAPIServer{messages: []message{{
Info: struct {
Role string `json:"role"`
}{Role: "assistant"},
Parts: []part{{Type: "text", Text: "a"}, {Type: "reasoning", Text: "x"}},
}}} }}}
c := fakeClient(t, f) c := fakeClient(t, f)
n, err := c.textCount(context.Background(), "sess-fake") msgs, err := c.Messages(context.Background(), "sess-fake")
if err != nil { if err != nil {
t.Fatalf("textCount err: %v", err) t.Fatalf("Messages err: %v", err)
} }
if n != 1 { if len(msgs) != 1 {
t.Errorf("textCount = %d, want 1 (одна text-часть в assistant)", n) t.Fatalf("len(msgs) = %d, want 1", len(msgs))
}
if !msgs[0].finished() {
t.Error("сообщение должно быть finished (Finish задан)")
}
}
func TestClient_Active(t *testing.T) {
f := &fakeAPIServer{active: true, sessionID: "sess-fake"}
c := fakeClient(t, f)
ok, err := c.Active(context.Background(), "sess-fake")
if err != nil {
t.Fatalf("Active err: %v", err)
}
if !ok {
t.Error("Active = false, want true")
}
ok, _ = c.Active(context.Background(), "sess-other")
if ok {
t.Error("Active(чужой) = true, want false")
}
}
func TestClient_Interrupt(t *testing.T) {
c := fakeClient(t, &fakeAPIServer{})
if err := c.Interrupt(context.Background(), "sess-fake"); err != nil {
t.Fatalf("Interrupt err: %v", err)
}
}
func Test_newestAssistant(t *testing.T) {
older := time.Now().Add(-time.Minute).UnixMilli()
newer := time.Now().UnixMilli()
msgs := []v2Message{
{ID: "a", Type: "assistant", Content: []v2Part{{Type: "text", Text: "x"}}, Time: v2Time{Created: &newer}},
{ID: "b", Type: "user", Time: v2Time{Created: &newer}},
{ID: "c", Type: "assistant", Content: []v2Part{{Type: "text", Text: "y"}}, Time: v2Time{Created: &older}},
}
cur, count := newestAssistant(msgs, older)
if cur == nil || cur.ID != "a" {
t.Errorf("newest = %v, want a", cur)
}
if count != 2 {
t.Errorf("count = %d, want 2", count)
}
texts := assistantText(msgs, older)
if len(texts) != 2 || texts[0] != "y" || texts[1] != "x" {
t.Errorf("assistantText order = %v, want [y x]", texts)
}
}
func Test_parseModelString(t *testing.T) {
m := parseModelString("tokentool/deepseek/deepseek-v4-flash-0731")
if m == nil || m.ProviderID != "tokentool" || m.ID != "deepseek/deepseek-v4-flash-0731" {
t.Errorf("parse = %+v, want tokentool/deepseek-v4-flash-0731", m)
}
if parseModelString("onlyprovider") != nil {
t.Error("parse без '/' должен вернуть nil")
}
if parseModelString("") != nil {
t.Error("parse пустой должен вернуть nil")
}
}
func TestClientErr_Unwrap(t *testing.T) {
ce := &ClientErr{Op: "prompt", Err: errors.New("boom")}
var target *ClientErr
if !errors.As(ce, &target) {
t.Fatal("expected *ClientErr")
} }
} }

176
internal/opencode/config.go Normal file
View File

@@ -0,0 +1,176 @@
package opencode
import (
"encoding/json"
"fmt"
"os"
"path/filepath"
"strings"
)
// Чтение top-level "model" из эффективного конфига opencode.
//
// Зачем: ratatoskr хардпинит модель в сессии (CreateSession), чтобы не зависеть
// от fallback-логики opencode. Если в конфиге модель не задана (или конфиг
// написан по старой v1-схеме — npm/options, которые v2 молча игнорирует),
// opencode сам выберет «дефолтную» модельную запись, и это может оказаться не
// той моделью. Поэтому мы явно логируем предупреждение (класс O5 WARN).
// opencodeConfigPath определяет путь к конфигу opencode, который видит
// serve-процесс этого пула (см. README): (1) явный OPENCODE_CONFIG из Server
// или окружения процесса, (2) OPENCODE_CONFIG_DIR / глобальный каталог
// ~/.config/opencode. Возвращает "" если ничего не найдено.
func opencodeConfigPath(cfgFile, cfgDir string) string {
// (1) явный файл конфига — Server.Config или env OPENCODE_CONFIG.
p := cfgFile
if p == "" {
p = os.Getenv("OPENCODE_CONFIG")
}
if p != "" {
if st, err := os.Stat(p); err == nil && !st.IsDir() {
return p
}
}
// (2) каталог конфигов.
dir := cfgDir
if dir == "" {
dir = os.Getenv("OPENCODE_CONFIG_DIR")
}
if dir == "" {
home, err := os.UserHomeDir()
if err != nil || home == "" {
return ""
}
dir = filepath.Join(home, ".config", "opencode")
if x := os.Getenv("XDG_CONFIG_HOME"); x != "" {
dir = filepath.Join(x, "opencode")
}
}
for _, name := range []string{"opencode.json", "opencode.jsonc"} {
cand := filepath.Join(dir, name)
if st, err := os.Stat(cand); err == nil && !st.IsDir() {
return cand
}
}
return ""
}
// ReadModelRef извлекает top-level "model" из конфига opencode и возвращает
// его как ModelRef. Модель не задана — вернёт (nil, nil); ошибка чтения/парсинга
// возвращается (вызывающий логирует warning и продолжает без хардпина).
func ReadModelRef(cfgFile, cfgDir string) (*ModelRef, error) {
path := opencodeConfigPath(cfgFile, cfgDir)
if path == "" {
return nil, nil
}
b, err := os.ReadFile(path)
if err != nil {
return nil, fmt.Errorf("config: читать %s: %w", path, err)
}
doc := struct {
Model json.RawMessage `json:"model"`
}{}
if err := json.Unmarshal(stripJSONC(b), &doc); err != nil {
return nil, fmt.Errorf("config: парсить %s: %w", path, err)
}
if len(doc.Model) == 0 || strings.TrimSpace(string(doc.Model)) == "null" {
return nil, nil
}
// "model" может быть строкой "provider/id" или объектом {providerID, id}.
var s string
if err := json.Unmarshal(doc.Model, &s); err == nil {
ref := parseModelString(s)
if ref == nil {
return nil, fmt.Errorf("config: некорректная model %q в %s (ожидается provider/id)", s, path)
}
return ref, nil
}
var ref ModelRef
if err := json.Unmarshal(doc.Model, &ref); err != nil {
return nil, fmt.Errorf("config: некорректная model в %s", path)
}
if ref.ProviderID == "" || ref.ID == "" {
return nil, fmt.Errorf("config: model без providerID/id в %s", path)
}
return &ref, nil
}
// parseModelString разбирает "provider/id" (как ModelV2.parse: провайдер — всё
// до первого '/', id — остаток). Возвращает nil при пустой/некорректной строке.
func parseModelString(s string) *ModelRef {
s = strings.TrimSpace(s)
if s == "" {
return nil
}
i := strings.IndexByte(s, '/')
if i <= 0 || i == len(s)-1 {
return nil
}
return &ModelRef{ProviderID: s[:i], ID: s[i+1:]}
}
// stripJSONC удаляет // и /* */ комментарии (вне строк), сохраняя позиции
// переводов строк, чтобы json.Unmarshal не споткнулся о trailing-комма.
func stripJSONC(b []byte) []byte {
out := make([]byte, 0, len(b))
inStr := false
esc := false
i := 0
for i < len(b) {
c := b[i]
if inStr {
out = append(out, c)
if esc {
esc = false
} else if c == '\\' {
esc = true
} else if c == '"' {
inStr = false
}
i++
continue
}
switch {
case c == '"':
inStr = true
out = append(out, c)
i++
case c == '/' && i+1 < len(b) && b[i+1] == '/':
for i < len(b) && b[i] != '\n' {
i++
}
if i < len(b) {
out = append(out, '\n')
i++
}
case c == '/' && i+1 < len(b) && b[i+1] == '*':
i += 2
for i+1 < len(b) && !(b[i] == '*' && b[i+1] == '/') {
i++
}
i += 2
default:
out = append(out, c)
i++
}
}
return dropTrailingCommas(out)
}
// dropTrailingCommas убирает запятые перед '}' / ']' (допускаются в JSONC).
func dropTrailingCommas(b []byte) []byte {
out := make([]byte, 0, len(b))
for i := 0; i < len(b); i++ {
if b[i] == ',' {
j := i + 1
for j < len(b) && (b[j] == ' ' || b[j] == '\t' || b[j] == '\n' || b[j] == '\r') {
j++
}
if j < len(b) && (b[j] == '}' || b[j] == ']') {
continue
}
}
out = append(out, b[i])
}
return out
}

View File

@@ -0,0 +1,85 @@
package opencode
import (
"os"
"path/filepath"
"testing"
)
func TestReadModelRef_String(t *testing.T) {
dir := t.TempDir()
path := filepath.Join(dir, "opencode.jsonc")
// конфиг с комментариями и trailing-запятыми (JSONC).
src := `{
// комментарий
"model": "tokentool/deepseek/deepseek-v4-flash-0731", /* и блочный */
"provider": {
"tokentool": {"api": {"type": "aisdk", "package": "@ai-sdk/openai-compatible", "url": "https://x"}},
},
}`
if err := os.WriteFile(path, []byte(src), 0o644); err != nil {
t.Fatalf("write: %v", err)
}
m, err := ReadModelRef(path, "")
if err != nil {
t.Fatalf("ReadModelRef err: %v", err)
}
if m == nil || m.ProviderID != "tokentool" || m.ID != "deepseek/deepseek-v4-flash-0731" {
t.Errorf("model = %+v, want tokentool/deepseek-v4-flash-0731", m)
}
}
func TestReadModelRef_Object(t *testing.T) {
dir := t.TempDir()
path := filepath.Join(dir, "opencode.json")
src := `{"model": {"providerID": "tokentool", "id": "deepseek/deepseek-v4-flash-0731"}}`
if err := os.WriteFile(path, []byte(src), 0o644); err != nil {
t.Fatalf("write: %v", err)
}
m, err := ReadModelRef(path, "")
if err != nil {
t.Fatalf("ReadModelRef err: %v", err)
}
if m == nil || m.ID != "deepseek/deepseek-v4-flash-0731" {
t.Errorf("model = %+v, want object-форма", m)
}
}
func TestReadModelRef_Missing(t *testing.T) {
dir := t.TempDir()
path := filepath.Join(dir, "opencode.json")
src := `{"provider": {}}`
if err := os.WriteFile(path, []byte(src), 0o644); err != nil {
t.Fatalf("write: %v", err)
}
m, err := ReadModelRef(path, "")
if err != nil {
t.Fatalf("ReadModelRef err: %v", err)
}
if m != nil {
t.Errorf("model = %+v, want nil (model не задан)", m)
}
}
func TestReadModelRef_NoFile(t *testing.T) {
dir := t.TempDir()
m, err := ReadModelRef(filepath.Join(dir, "nope.json"), dir)
if err != nil {
t.Fatalf("ReadModelRef err: %v", err)
}
if m != nil {
t.Errorf("model = %+v, want nil", m)
}
}
func TestReadModelRef_Bad(t *testing.T) {
dir := t.TempDir()
path := filepath.Join(dir, "opencode.json")
src := `{"model": 12345}`
if err := os.WriteFile(path, []byte(src), 0o644); err != nil {
t.Fatalf("write: %v", err)
}
if _, err := ReadModelRef(path, ""); err == nil {
t.Error("ReadModelRef должен упасть на некорректной model")
}
}

View File

@@ -6,7 +6,7 @@ import (
"fmt" "fmt"
"io" "io"
"os" "os"
"sync" "strings"
"time" "time"
) )
@@ -20,11 +20,13 @@ type Result struct {
SessionID string SessionID string
} }
// Runner — запуск opencode-субагентов через HTTP API serve. // Runner — запуск opencode-субагентов через v2 HTTP API serve.
// //
// Полный переход на API: Runner ходит к opencode serve через Pool→Client // Runner ходит к opencode serve через Pool→Client (пути /api/*, см. README,
// (нет spawn-модели, нет NDJSON). Агент идёт в сервер пула для своего каталога // минимальная версия opencode). Промпт отправляется неблокирующе (durable
// (в нём запущен serve → он его project). // admit), вердикт собирается поллингом новых assistant-сообщений; завершение
// ответа определяется по схеме «сессия больше не в активных дренажах» + финальное
// assistant-сообщение.
type Runner struct { type Runner struct {
Pool *Pool // пул serve-серверов (обязательный) Pool *Pool // пул serve-серверов (обязательный)
IdleTimeout time.Duration IdleTimeout time.Duration
@@ -73,49 +75,64 @@ func (r *Runner) Run(ctx context.Context, prompt, cwd, agent, sessionID string)
} }
c := &Client{BaseURL: srv.Addr(), Password: srv.Password, Debug: r.Debug} c := &Client{BaseURL: srv.Addr(), Password: srv.Password, Debug: r.Debug}
// Модель по умолчанию из конфига opencode — хардпиним её в сессии, чтобы
// не зависеть от fallback-логики opencode (класс O5 WARN: если модель не
// считывается/не задана — предупреждаем и работаем без явного указания).
model, mErr := ReadModelRef(srv.Config, srv.ConfigDir)
if mErr != nil {
r.logf("WARN opencode: не удалось прочитать model из конфига: %v", mErr)
} else if model == nil {
r.logf("WARN opencode: в конфиге opencode не задан top-level model — модель не хардпинится (риск fallback)")
} else {
r.logf("opencode(%s) model=%s", agent, model)
}
// Сессия: заданная (resume) или новая. // Сессия: заданная (resume) или новая.
sid := sessionID sid := sessionID
if sid == "" { if sid == "" {
sid, err = c.CreateSession(ctx, "ratatoskr-"+agent) sid, err = c.CreateSession(ctx, model)
if err != nil { if err != nil {
return nil, fmt.Errorf("opencode: create session: %w", err) return nil, fmt.Errorf("opencode: create session: %w", err)
} }
r.logf("opencode(%s) session=%s на %s", agent, sid, srv.Addr()) r.logf("opencode(%s) session=%s на %s", agent, sid, srv.Addr())
} }
// Отправляем промпт (блокирующий Send в горутине; вердикт придёт из него), return r.awaitVerdict(ctx, c, model, sid, agent, prompt)
// параллельно поллим прогресс и контролируем idle/hard таймауты.
return r.awaitVerdict(ctx, c, sid, agent, prompt)
} }
// awaitVerdict запускает блокирующий Send и параллельно поллит прогресс // settlePolls — сколько подряд опросов должно подтвердить завершение ответа,
// (рост числа text-частей = агент жив, сбрасывает idle). Возвращается вердикт // прежде чем считать вердикт финальным (устойчивость к гонке между удалением
// из ответа Send, либо rc=-1 при idle/hard таймауте (тогда Abort + отмена ctx). // сессии из активных дренажей и финализацией последнего сообщения).
func (r *Runner) awaitVerdict(ctx context.Context, c *Client, sid, agent, prompt string) (*Result, error) { const settlePolls = 2
sendCtx, cancel := context.WithCancel(ctx)
defer cancel()
type sendOut struct {
vd string
err error
}
sendCh := make(chan sendOut, 1)
go func() {
vd, err := c.Send(sendCtx, sid, prompt)
sendCh <- sendOut{vd: vd, err: err}
}()
// Прогресс = сумма text-частей во всех assistant-сообщениях сессии. Рост // awaitVerdict отправляет промпт (неблокирующе) и поллит новые assistant-сообщения,
// сбрасывает idle-таймер (LLM стримит = жив). // контролируя idle/hard таймауты. Завершение: сессия ушла из активных дренажей
var mu sync.Mutex // И есть новое завершённое assistant-сообщение, стабильное в течение settlePolls
lastCount := -1 // опросов. Возвращает вердикт (текст text-партов), либо rc=-1 при таймауте.
func (r *Runner) awaitVerdict(ctx context.Context, c *Client, model *ModelRef, sid, agent, prompt string) (*Result, error) {
// admit промпта; граница «новых» сообщений — время создания user-сообщения.
admittedAt := time.Now().UnixMilli()
adm, err := c.Prompt(ctx, sid, prompt)
if err != nil {
return nil, err
}
if adm != nil && adm.TimeCreated > 0 {
admittedAt = adm.TimeCreated
}
// Прогресс = число контент-партов + суммарная длина их текста в новых
// assistant-сообщениях (progressOf). Рост сбрасывает idle-таймер: LLM
// стримит (даже в один растущий text-парт) или думает (reasoning) = жив.
lastParts, lastTextLen := -1, -1
lastProgress := time.Now() lastProgress := time.Now()
launch := time.Now() launch := time.Now()
doneSeen, emptySeen := 0, 0
abortAnd := func(rc int, why string) (*Result, error) { abortAnd := func(rc int, why string) (*Result, error) {
if err := c.Abort(ctx, sid); err != nil { if err := c.Interrupt(ctx, sid); err != nil {
r.logf("opencode(%s) abort %s: %v", agent, why, err) r.logf("opencode(%s) interrupt %s: %v", agent, why, err)
} }
cancel()
return &Result{RC: rc, Stdout: "", SessionID: sid}, nil return &Result{RC: rc, Stdout: "", SessionID: sid}, nil
} }
@@ -125,43 +142,85 @@ func (r *Runner) awaitVerdict(ctx context.Context, c *Client, sid, agent, prompt
return abortAnd(-1, "ctx") return abortAnd(-1, "ctx")
} }
count, _ := c.textCount(ctx, sid) msgs, err := c.Messages(ctx, sid)
mu.Lock() if err != nil {
if count != lastCount { if ctx.Err() != nil {
lastProgress = time.Now() return abortAnd(-1, "ctx")
lastCount = count }
var ce *ClientErr
if errors.As(err, &ce) && ce.Op == "connect" {
return nil, fmt.Errorf("opencode: %w", err)
}
return nil, err
}
active, err := c.Active(ctx, sid)
if err != nil {
if ctx.Err() != nil {
return abortAnd(-1, "ctx")
}
var ce *ClientErr
if errors.As(err, &ce) && ce.Op == "connect" {
return nil, fmt.Errorf("opencode: %w", err)
}
return nil, err
} }
mu.Unlock()
cur, _ := newestAssistant(msgs, admittedAt)
parts, textLen := progressOf(msgs, admittedAt)
if parts != lastParts || textLen != lastTextLen {
lastProgress = time.Now()
lastParts, lastTextLen = parts, textLen
}
now := time.Now() now := time.Now()
if now.Sub(lastProgress) > r.IdleTimeout { if now.Sub(lastProgress) > r.IdleTimeout {
r.logf("opencode(%s) idle %.0fs — abort", agent, r.IdleTimeout.Seconds()) r.logf("opencode(%s) idle %.0fs — abort", agent, r.IdleTimeout.Seconds())
return abortAnd(-1, "idle") return abortAnd(-1, "idle")
} }
// hard — общий бюджет от старта запуска.
if now.Sub(launch) > r.HardTimeout { if now.Sub(launch) > r.HardTimeout {
r.logf("opencode(%s) hard timeout %.0fs — abort", agent, r.HardTimeout.Seconds()) r.logf("opencode(%s) hard timeout %.0fs — abort", agent, r.HardTimeout.Seconds())
return abortAnd(-1, "hard") return abortAnd(-1, "hard")
} }
switch {
case !active && cur != nil && cur.finished():
// ответ закончен — ждём стабильности, затем собираем вердикт
doneSeen++
emptySeen = 0
if doneSeen >= settlePolls {
return r.verdict(model, cur, msgs, admittedAt, sid)
}
case !active && cur == nil:
// сессия завершилась, но нового assistant-сообщения так и нет
emptySeen++
if emptySeen >= settlePolls {
return nil, &ClientErr{Op: "prompt", Err: errors.New("агент не выдал ответ (сессия пуста)")}
}
default:
doneSeen, emptySeen = 0, 0
}
select { select {
case out := <-sendCh:
// Send завершился. Ошибка — connect (сервер недоступен) и ctx жив →
// фатально, не таймаут. Если ctx уже отменён — это обрыв, а не ошибка.
if out.err != nil {
var ce *ClientErr
if errors.As(out.err, &ce) && ce.Op == "connect" && ctx.Err() == nil {
return nil, fmt.Errorf("opencode: %w", out.err)
}
if ctx.Err() != nil {
return abortAnd(-1, "ctx")
}
return nil, out.err
}
r.logf("opencode(%s) вердикт готов (%d байт)", agent, len(out.vd))
return &Result{RC: 0, Stdout: out.vd, SessionID: sid}, nil
case <-time.After(r.PollInterval): case <-time.After(r.PollInterval):
case <-ctx.Done(): case <-ctx.Done():
} }
} }
} }
// verdict собирает финальный результат из новых assistant-сообщений.
// Проверяет фактическую модель ответа и логирует warning при расхождении
// с ожидаемой (устойчивость к «не той» модели — класс O5 WARN).
func (r *Runner) verdict(model *ModelRef, cur *v2Message, msgs []v2Message, since int64, sid string) (*Result, error) {
if model != nil && cur.Model != nil && (model.ProviderID != cur.Model.ProviderID || model.ID != cur.Model.ID) {
r.logf("WARN opencode: сессия %s отвечала моделью %s, а не ожидаемой %s — проверь providers в конфиге (v2-схема: provider.api / request, а не npm/options)", sid, cur.Model, model)
}
if cur.Error != nil && cur.Error.Message != "" {
return nil, &ClientErr{Op: "prompt", Err: errors.New(cur.Error.Message)}
}
texts := assistantText(msgs, since)
if len(texts) == 0 {
return nil, &ClientErr{Op: "prompt", Err: errors.New("нет text-части в ответе")}
}
vd := stripFence(strings.Join(texts, "\n"))
r.logf("opencode вердикт готов (%d байт)", len(vd))
return &Result{RC: 0, Stdout: vd, SessionID: sid}, nil
}

View File

@@ -2,6 +2,7 @@ package opencode
import ( import (
"context" "context"
"io"
"net/http/httptest" "net/http/httptest"
"testing" "testing"
"time" "time"
@@ -9,6 +10,8 @@ import (
// fakePool создаёт Pool, в котором уже «живёт» сервер для каталога (без spawn): // fakePool создаёт Pool, в котором уже «живёт» сервер для каталога (без spawn):
// Server{URL: fake.URL}, поэтому Runner ходит по HTTP на фейк-API. // Server{URL: fake.URL}, поэтому Runner ходит по HTTP на фейк-API.
// XDG_CONFIG_HOME уводится во временный каталог, чтобы ReadModelRef не читал
// реальный пользовательский конфиг opencode (детерминизм тестов).
func fakePool(t *testing.T, f *fakeAPIServer, dir string) (*Pool, *Client) { func fakePool(t *testing.T, f *fakeAPIServer, dir string) (*Pool, *Client) {
t.Helper() t.Helper()
ts := httptestURL(t, f) ts := httptestURL(t, f)
@@ -28,13 +31,12 @@ func httptestURL(t *testing.T, f *fakeAPIServer) string {
} }
func TestRun_Success(t *testing.T) { func TestRun_Success(t *testing.T) {
t.Setenv("XDG_CONFIG_HOME", t.TempDir())
dir := t.TempDir() dir := t.TempDir()
f := &fakeAPIServer{ f := &fakeAPIServer{verdictText: "done"}
verdictParts: []part{{Type: "text", Text: "done"}},
}
p, _ := fakePool(t, f, dir) p, _ := fakePool(t, f, dir)
r := &Runner{Pool: p, PollInterval: 5 * time.Millisecond} r := &Runner{Pool: p, PollInterval: 5 * time.Millisecond, Stdout: io.Discard}
res, err := r.Run(context.Background(), "task", dir, "dev", "") res, err := r.Run(context.Background(), "task", dir, "dev", "")
if err != nil { if err != nil {
t.Fatalf("Run err: %v", err) t.Fatalf("Run err: %v", err)
@@ -51,13 +53,14 @@ func TestRun_Success(t *testing.T) {
} }
func TestRun_IdleTimeout(t *testing.T) { func TestRun_IdleTimeout(t *testing.T) {
t.Setenv("XDG_CONFIG_HOME", t.TempDir())
dir := t.TempDir() dir := t.TempDir()
// prompt блокируется (агент «завис»), прогресс не растёт → idle abort // агент «завис»: active=true, прогресс не растёт → idle abort
f := &fakeAPIServer{blockPrompt: true} f := &fakeAPIServer{blockPrompt: true}
p, _ := fakePool(t, f, dir) p, _ := fakePool(t, f, dir)
r := &Runner{Pool: p, IdleTimeout: 30 * time.Millisecond, r := &Runner{Pool: p, IdleTimeout: 30 * time.Millisecond,
PollInterval: 5 * time.Millisecond} PollInterval: 5 * time.Millisecond, Stdout: io.Discard}
res, err := r.Run(context.Background(), "task", dir, "dev", "") res, err := r.Run(context.Background(), "task", dir, "dev", "")
if err != nil { if err != nil {
t.Fatalf("Run err: %v", err) t.Fatalf("Run err: %v", err)
@@ -67,14 +70,55 @@ func TestRun_IdleTimeout(t *testing.T) {
} }
} }
func TestRun_StreamingGrowth(t *testing.T) {
t.Setenv("XDG_CONFIG_HOME", t.TempDir())
dir := t.TempDir()
// стрим: один text-парт растёт с каждым опросом (дольше, чем idle timeout),
// но модель жива → idle НЕ должен сработать.
f := &fakeAPIServer{streamGrow: true, streamPolls: 30}
p, _ := fakePool(t, f, dir)
r := &Runner{Pool: p, IdleTimeout: 60 * time.Millisecond,
PollInterval: 5 * time.Millisecond, Stdout: io.Discard}
res, err := r.Run(context.Background(), "task", dir, "dev", "")
if err != nil {
t.Fatalf("Run err: %v", err)
}
if res.RC != 0 {
t.Errorf("RC = %d, want 0 (растущий стрим не должен считаться hung)", res.RC)
}
if !contains(res.Stdout, "done-stream") {
t.Errorf("Stdout = %q, want contain done-stream", res.Stdout)
}
}
func TestRun_ReasoningGrowth(t *testing.T) {
t.Setenv("XDG_CONFIG_HOME", t.TempDir())
dir := t.TempDir()
// стрим: растёт только reasoning-парт (текста нет) — тоже живая активность.
f := &fakeAPIServer{streamGrow: true, streamReasoning: true, streamPolls: 30}
p, _ := fakePool(t, f, dir)
r := &Runner{Pool: p, IdleTimeout: 60 * time.Millisecond,
PollInterval: 5 * time.Millisecond, Stdout: io.Discard}
res, err := r.Run(context.Background(), "task", dir, "dev", "")
if err != nil {
t.Fatalf("Run err: %v", err)
}
if res.RC != 0 {
t.Errorf("RC = %d, want 0 (растущий reasoning не должен считаться hung)", res.RC)
}
}
func TestRun_ContextCancel(t *testing.T) { func TestRun_ContextCancel(t *testing.T) {
t.Setenv("XDG_CONFIG_HOME", t.TempDir())
dir := t.TempDir() dir := t.TempDir()
f := &fakeAPIServer{blockPrompt: true} f := &fakeAPIServer{blockPrompt: true}
p, _ := fakePool(t, f, dir) p, _ := fakePool(t, f, dir)
ctx, cancel := context.WithCancel(context.Background()) ctx, cancel := context.WithCancel(context.Background())
r := &Runner{Pool: p, IdleTimeout: time.Minute, HardTimeout: time.Minute, r := &Runner{Pool: p, IdleTimeout: time.Minute, HardTimeout: time.Minute,
PollInterval: 5 * time.Millisecond} PollInterval: 5 * time.Millisecond, Stdout: io.Discard}
done := make(chan *Result, 1) done := make(chan *Result, 1)
errCh := make(chan error, 1) errCh := make(chan error, 1)
go func() { go func() {

View File

@@ -143,9 +143,15 @@ func (s *Server) serveCmd(ctx context.Context) *exec.Cmd {
return cmd return cmd
} }
// waitHealthy опрашивает /global/health сервера до первого успеха или Connect. // MinVersion — минимальная версия opencode, с которой работает интеграция.
// v2 HTTP API (префикс /api/*) присутствует в сборках dev / >=1.18.18.
// Более старые бинари отвечают на /global/health и НЕ подходят.
const MinVersion = "1.18.18"
// waitHealthy опрашивает /api/health сервера до первого успеха или Connect.
// Возвращает nil, как только сервер ответил {healthy:true} (или 200/401 — сервер // Возвращает nil, как только сервер ответил {healthy:true} (или 200/401 — сервер
// жив, но может требовать авторизации). // жив, но может требовать авторизации). При неудаче — ошибка с подсказкой про
// минимальную версию opencode (класс O1: старый бинарь не знает v2-путей).
func (s *Server) waitHealthy(ctx context.Context, addr string) error { func (s *Server) waitHealthy(ctx context.Context, addr string) error {
deadline := time.Now().Add(60 * time.Second) deadline := time.Now().Add(60 * time.Second)
poll := s.PollInterval poll := s.PollInterval
@@ -154,7 +160,7 @@ func (s *Server) waitHealthy(ctx context.Context, addr string) error {
return nil return nil
} }
if time.Now().After(deadline) { if time.Now().After(deadline) {
return fmt.Errorf("opencode serve %s: не стал доступным (healthcheck)", addr) return fmt.Errorf("opencode serve %s: не стал доступным (v2 healthcheck). Нужен opencode >= %s (v2 HTTP API /api/*), а не старый бинарь", addr, MinVersion)
} }
select { select {
case <-ctx.Done(): case <-ctx.Done():
@@ -167,7 +173,7 @@ func (s *Server) waitHealthy(ctx context.Context, addr string) error {
// healthGET делает GET на адрес и возвращает true, если сервер ответил. // healthGET делает GET на адрес и возвращает true, если сервер ответил.
// 401 (basic auth требуется) тоже считается «жив» — сервер доступен. // 401 (basic auth требуется) тоже считается «жив» — сервер доступен.
func (s *Server) healthGET(ctx context.Context, addr string) bool { func (s *Server) healthGET(ctx context.Context, addr string) bool {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, addr+"/global/health", nil) req, err := http.NewRequestWithContext(ctx, http.MethodGet, addr+"/api/health", nil)
if err != nil { if err != nil {
return false return false
} }

View File

@@ -8,15 +8,25 @@ import (
"os" "os"
"os/exec" "os/exec"
"path/filepath" "path/filepath"
"runtime"
"strings" "strings"
"testing" "testing"
"time" "time"
) )
// fakeServeBin создаёт скрипт, имитирующий opencode serve: просто держит // fakeServeBin создаёт скрипт, имитирующий opencode serve: просто держит
// процесс живым (sleep), чтобы супервайзер мог им владеть и убивать его. // процесс живым (sleep/ping), чтобы супервайзер мог им владеть и убивать его.
// На Windows используется .cmd (с #!/bin/sh нельзя — он не исполняется).
func fakeServeBin(t *testing.T, workdir string) string { func fakeServeBin(t *testing.T, workdir string) string {
t.Helper() t.Helper()
if runtime.GOOS == "windows" {
bin := filepath.Join(workdir, "opencode-serve.cmd")
script := "@echo off\r\necho fake serve started\r\nping -n 300 127.0.0.1 >nul\r\n"
if err := os.WriteFile(bin, []byte(script), 0o755); err != nil {
t.Fatalf("write fake serve bin: %v", err)
}
return bin
}
bin := filepath.Join(workdir, "opencode-serve") bin := filepath.Join(workdir, "opencode-serve")
script := `#!/bin/sh script := `#!/bin/sh
echo "fake serve started" echo "fake serve started"

View File

@@ -1,66 +1,43 @@
package storage package storage
import "encoding/json" import (
"encoding/json"
// Status — статус задачи (state machine). "github.com/kamelion/ratatoskr-go/internal/model"
type Status string )
// Status — алиас доменного статуса задачи.
//
// Совместимый мост: весь внешний код продолжает использовать storage.Status
// (например «storage.StatusRunning»), но единый источник истины — model.Status.
type Status = model.Status
// Статусы задачи — re-export из model.
const ( const (
StatusDraft Status = "draft" // только что создана StatusDraft = model.StatusDraft
StatusCollecting Status = "collecting" // аналитик собирает детали StatusCollecting = model.StatusCollecting
StatusReady Status = "ready" // черновик готов, ждёт одобрения пользователя StatusReady = model.StatusReady
StatusApproved Status = "approved" // пользователь одобрил («создавай») — воркер берёт в работу StatusApproved = model.StatusApproved
StatusRunning Status = "running" // opencode работает StatusRunning = model.StatusRunning
StatusSuccess Status = "success" // задача выполнена StatusSuccess = model.StatusSuccess
StatusFailed Status = "failed" // ошибка выполнения StatusFailed = model.StatusFailed
StatusTimeout Status = "timeout" // таймаут opencode StatusTimeout = model.StatusTimeout
StatusCancelled Status = "cancelled" // отменена пользователем StatusCancelled = model.StatusCancelled
StatusAborted Status = "aborted" // сбой сбора, черновик выброшен StatusAborted = model.StatusAborted
StatusClosed Status = "closed" // закрыта вручную StatusClosed = model.StatusClosed
) )
// AllStatuses — все возможные статусы для валидации. // AllStatuses — все возможные статусы для валидации.
var AllStatuses = []Status{ var AllStatuses = model.AllStatuses
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. // IsValidTransition проверяет, допустим ли переход from → to.
func IsValidTransition(from, to Status) bool { func IsValidTransition(from, to Status) bool {
allowed, ok := validTransitions[from] return model.IsValidTransition(from, to)
if !ok {
return false
}
for _, s := range allowed {
if s == to {
return true
}
}
return false
} }
// IsTerminal возвращает true, если статус терминальный. // IsTerminal возвращает true, если статус терминальный.
func IsTerminal(s Status) bool { func IsTerminal(s Status) bool {
return s == StatusSuccess || s == StatusCancelled || return model.IsTerminal(s)
s == StatusAborted || s == StatusClosed
} }
// Task — запись задачи в БД. // Task — запись задачи в БД.
@@ -111,14 +88,14 @@ func (t *Task) SetReposFromDB(repos string) {
_ = json.Unmarshal([]byte(repos), &t.Repos) _ = json.Unmarshal([]byte(repos), &t.Repos)
} }
// Trace — запись трассировки выполнения. // TraceStatus — алиас доменного статуса трассировки.
type TraceStatus string type TraceStatus = model.TraceStatus
const ( const (
TraceRunning TraceStatus = "running" TraceRunning TraceStatus = model.TraceRunning
TraceSuccess TraceStatus = "success" TraceSuccess TraceStatus = model.TraceSuccess
TraceFailed TraceStatus = "failed" TraceFailed TraceStatus = model.TraceFailed
TraceTimeout TraceStatus = "timeout" TraceTimeout TraceStatus = model.TraceTimeout
) )
// Trace — лог одного субагента. // Trace — лог одного субагента.

46
internal/ui/commands.go Normal file
View File

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

52
internal/ui/controller.go Normal file
View File

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

View File

@@ -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() // не должно падать
}

View File

@@ -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:])
}

73
internal/ui/store.go Normal file
View File

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

107
internal/ui/store_test.go Normal file
View File

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

40
internal/ui/view.go Normal file
View File

@@ -0,0 +1,40 @@
// Package ui — контроллер и представления десктопного интерфейса (Fyne).
//
// Слой обмена с ядром: однонаправленный поток (спец 12.212.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) {}

68
internal/ui/window.go Normal file
View File

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

View File

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

View File

@@ -10,6 +10,7 @@ import (
"strings" "strings"
"time" "time"
"github.com/kamelion/ratatoskr-go/internal/events"
"github.com/kamelion/ratatoskr-go/internal/opencode" "github.com/kamelion/ratatoskr-go/internal/opencode"
"github.com/kamelion/ratatoskr-go/internal/storage" "github.com/kamelion/ratatoskr-go/internal/storage"
) )
@@ -56,6 +57,27 @@ type Worker struct {
// подменяемый poll для тестов // подменяемый poll для тестов
pollFn PollTaskFunc 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 оборачивает контекст запуска субагента, привязывая живое наблюдение // runCtx оборачивает контекст запуска субагента, привязывая живое наблюдение
@@ -192,8 +214,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
} }
// 2. ставим running (после валидации — чтобы плохие имена не жгли состояние) // 2. ставим running (после валидации — чтобы плохие имена не жгли состояние)
task.Status = storage.StatusRunning if err := w.setStatus(ctx, task, storage.StatusRunning); err != nil {
if err := w.Store.UpdateTask(ctx, task); err != nil {
return fmt.Errorf("%w: set running: %v", ErrUpdate, err) return fmt.Errorf("%w: set running: %v", ErrUpdate, err)
} }
w.notifyStatus(ctx, task, storage.StatusRunning) 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 0:
// продолжаем на ревью // продолжаем на ревью
case -1: case -1:
task.Status = storage.StatusTimeout if e := w.setStatus(ctx, task, storage.StatusTimeout); e != nil {
if e := w.Store.UpdateTask(ctx, task); e != nil {
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e) return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
} }
w.notifyStatus(ctx, task, storage.StatusTimeout) w.notifyStatus(ctx, task, storage.StatusTimeout)
w.finalizeTrace(ctx, traceID, storage.TraceTimeout, output) w.finalizeTrace(ctx, traceID, storage.TraceTimeout, output)
return nil return nil
default: default:
task.Status = storage.StatusFailed if e := w.setStatus(ctx, task, storage.StatusFailed); e != nil {
if e := w.Store.UpdateTask(ctx, task); e != nil {
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e) return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
} }
w.notifyStatus(ctx, task, storage.StatusFailed) 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 { if verdict == nil {
// невалидный JSON даже после retry → failed с объяснением. // невалидный JSON даже после retry → failed с объяснением.
task.Status = storage.StatusFailed
explain := "reviewer вернул невалидный/пустой вердикт (даже после повтора)." 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) return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
} }
w.notifyStatus(ctx, task, storage.StatusFailed) w.notifyStatus(ctx, task, storage.StatusFailed)
@@ -325,8 +343,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
w.failTask(ctx, task) w.failTask(ctx, task)
return pErr return pErr
} }
task.Status = storage.StatusSuccess if e := w.setStatus(ctx, task, storage.StatusSuccess); e != nil {
if e := w.Store.UpdateTask(ctx, task); e != nil {
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e) return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
} }
w.notifyStatus(ctx, task, storage.StatusSuccess) w.notifyStatus(ctx, task, storage.StatusSuccess)
@@ -341,8 +358,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
} }
// Лимит исчерпан → failed с объяснением. // Лимит исчерпан → failed с объяснением.
task.Status = storage.StatusFailed if e := w.setStatus(ctx, task, storage.StatusFailed); e != nil {
if e := w.Store.UpdateTask(ctx, task); e != nil {
return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e) return fmt.Errorf("%w: set %s: %v", ErrUpdate, task.Status, e)
} }
w.notify(ctx, task, fmt.Sprintf("Задача #%d: failed — ревью не пройдено за %d итераций", task.ID, maxReviewIterations)) 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 и уведомляет владельца. // failTask помечает задачу failed и уведомляет владельца.
func (w *Worker) failTask(ctx context.Context, task *storage.Task) { func (w *Worker) failTask(ctx context.Context, task *storage.Task) {
task.Status = storage.StatusFailed if e := w.setStatus(ctx, task, storage.StatusFailed); e != nil {
if e := w.Store.UpdateTask(ctx, task); e != nil {
log.Printf("worker: task %d: set failed: %v", task.ID, e) log.Printf("worker: task %d: set failed: %v", task.ID, e)
} }
w.notifyStatus(ctx, task, storage.StatusFailed) w.notifyStatus(ctx, task, storage.StatusFailed)

View File

@@ -4,7 +4,6 @@ import (
"context" "context"
"encoding/json" "encoding/json"
"errors" "errors"
"fmt"
"os" "os"
"os/exec" "os/exec"
"path/filepath" "path/filepath"
@@ -809,13 +808,41 @@ func TestWorkerStartStop(t *testing.T) {
w.Stop() w.Stop()
} }
// waitTaskStatus ждёт, пока задача достигнет статуса want. На Windows git-операции
// воркера заметно медленнее, чем на Linux, поэтому проверки в тестах не могут
// полагаться на фиксированные sleep'ы — только на polling до целевого статуса.
func waitTaskStatus(t *testing.T, ctx context.Context, s *storage.Storage, id int64, want storage.Status) {
t.Helper()
deadline := time.Now().Add(30 * time.Second)
for {
task, err := s.GetTask(ctx, id)
if err != nil {
t.Fatalf("get task %d: %v", id, err)
}
if task.Status == want {
return
}
switch task.Status {
case storage.StatusFailed, storage.StatusTimeout:
t.Fatalf("task %d: status %q, want %q", id, task.Status, want)
}
if time.Now().After(deadline) {
t.Fatalf("task %d: таймаут ожидания %q, последний статус %q", id, want, task.Status)
}
select {
case <-ctx.Done():
t.Fatalf("task %d: ctx done: %v", id, ctx.Err())
case <-time.After(50 * time.Millisecond):
}
}
}
func TestWorkerSemaphore(t *testing.T) { func TestWorkerSemaphore(t *testing.T) {
s := setupWorkerDB(t) s := setupWorkerDB(t)
// создаём 2 ready-задачи // создаём 2 ready-задачи
for i := 0; i < 2; i++ { task1 := createReadyTask(t, s, "task-0")
createReadyTask(t, s, fmt.Sprintf("task-%d", i)) task2 := createReadyTask(t, s, "task-1")
}
w := &Worker{ w := &Worker{
Store: s, Store: s,
@@ -829,12 +856,12 @@ func TestWorkerSemaphore(t *testing.T) {
w.sem = make(chan struct{}, 1) w.sem = make(chan struct{}, 1)
w.sem <- struct{}{} w.sem <- struct{}{}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel() defer cancel()
// первый poll — запустит 1 задачу (макс. 1) // первый poll — запустит 1 задачу (макс. 1)
w.pollAndDispatch(ctx) w.pollAndDispatch(ctx)
time.Sleep(200 * time.Millisecond) waitTaskStatus(t, ctx, s, task1.ID, storage.StatusSuccess)
// 1 должна быть success, 1 — всё ещё approved // 1 должна быть success, 1 — всё ещё approved
success, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess}) success, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess})
@@ -848,7 +875,7 @@ func TestWorkerSemaphore(t *testing.T) {
// первая завершилась и вернула токен в сем — можем диспатчить вторую // первая завершилась и вернула токен в сем — можем диспатчить вторую
w.pollAndDispatch(ctx) w.pollAndDispatch(ctx)
time.Sleep(200 * time.Millisecond) waitTaskStatus(t, ctx, s, task2.ID, storage.StatusSuccess)
success, _ = s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess}) success, _ = s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess})
if len(success) != 2 { if len(success) != 2 {

View File

@@ -0,0 +1,109 @@
<#
build-publish-ui.ps1 — локальная сборка Windows-бинаря с UI (Fyne) и публикация в Gitea Packages.
Почему локально: Fyne-окно требует cgo + C-компилятор (MinGW/gcc), которых нет в CI.
CI собирает ТОЛЬКО headless-linux (см. .gitea/workflows/ci.yaml). Windows-бинарь
с графическим окном собирается на Windows-машине Камиля и публикуется скриптом ниже
в ту же версию commit-<sha7>, что и CI, с теми же companion-файлами (.version/.sha256),
чтобы автообновление (U4 checksum, ResolveLatest) работало без изменений.
Требования:
- Go с доступным C-компилятором (CGO_ENABLED=1). WinLibs gcc: см. память
`mem:toolchain/cgo-winlibs-gcc` — bin нужно добавить в PATH перед запуском.
- Файл scripts/.gitea-creds (git-ignored, см. ниже) с токенами.
Формат (по одной `ключ=значение` на строку):
GITEA_TOKEN=write:packages-токен
UPDATE_TOKEN=read:package-токен
GIT_MAIN_URL=http://gitea.hal9000.home
- Ветка уже запушена в `main` (VERSION = commit-<sha7> от текущего HEAD).
Пример:
powershell -ExecutionPolicy Bypass -File scripts/build-publish-ui.ps1
#>
$ErrorActionPreference = "Stop"
$ScriptDir = Split-Path -Parent $MyInvocation.MyCommand.Path
$RepoRoot = Split-Path -Parent $ScriptDir
$CredsFile = Join-Path $ScriptDir ".gitea-creds"
$Package = "ratatoskr"
$Owner = "kamelion"
$Filename = "ratatoskr-windows-amd64.exe"
Set-Location $RepoRoot
# --- 0. Проверка C-тулчейна ---
if ($env:CGO_ENABLED -eq "0") { Write-Error "CGO_ENABLED=0 — нужен C-компилятор (MinGW). Уберите его из env." }
if (-not (Get-Command gcc -ErrorAction SilentlyContinue)) {
Write-Host "gcc не найден в PATH. Пример (WinLibs):"
Write-Host ' $p = "<...>\mingw64\bin"; $env:Path = "$p;$env:Path"'
Write-Error "C-компилятор gcc не найден."
}
Write-Host "[ok] gcc: $((gcc --version | Select-Object -First 1))"
# --- 1. Чтение секретов ---
if (-not (Test-Path $CredsFile)) {
Write-Error "Нет файла секретов: $CredsFile`nФормат: GITEA_TOKEN=... / UPDATE_TOKEN=... / GIT_MAIN_URL=..."
}
$creds = @{}
Get-Content $CredsFile | ForEach-Object {
$line = $_.Trim()
if ($line -and -not $line.StartsWith("#")) {
$kv = $line.Split("=", 2)
if ($kv.Count -eq 2) { $creds[$kv[0].Trim()] = $kv[1].Trim() }
}
}
$giteaToken = $creds["GITEA_TOKEN"]
$updateToken = $creds["UPDATE_TOKEN"]
$giteaUrl = $creds["GIT_MAIN_URL"]
if (-not $giteaToken) { Write-Error "В $CredsFile нет GITEA_TOKEN (нужен write:packages)." }
if (-not $giteaUrl) { Write-Error "В $CredsFile нет GIT_MAIN_URL." }
# --- 2. Версия = commit-<sha7> текущего HEAD (совпадает с именованием CI) ---
$sha = (git rev-parse --short HEAD).Trim()
if (-not $sha) { Write-Error "git rev-parse --short HEAD не дал хэш." }
$Version = "commit-$sha"
Write-Host "[build] VERSION=$Version FILENAME=$Filename"
# --- 3. Сборка (cgo + UI) ---
$env:CGO_ENABLED = "1"
$out = Join-Path $RepoRoot $Filename
$ldflags = "-s -w -X main.version=${Version} -X main.updateToken=${updateToken}"
Write-Host "[build] go build -ldflags=... -o $Filename"
$goArgs = @("build", "-ldflags=$ldflags", "-o", $out, "./cmd/ratatoskr/")
go $goArgs
if ($LASTEXITCODE -ne 0) { Write-Error "go build завершился с кодом $LASTEXITCODE" }
# --- 4. Companion-файлы (.version, .sha256) ---
$sha256Hex = (Get-FileHash -Algorithm SHA256 $out).Hash.ToLower()
$verFile = "$out.version"
$sumFile = "$out.sha256"
Set-Content -Path $verFile -Value $Version -NoNewline -Encoding ascii
Set-Content -Path $sumFile -Value $sha256Hex -NoNewline -Encoding ascii
Write-Host "[ok] bin: $out ($([math]::Round((Get-Item $out).Length/1MB,1)) MB)"
Write-Host "[ok] sha256: $sha256Hex"
# --- 5. Публикация в Gitea Packages ---
$api = "$giteaUrl/api/packages/$Owner/generic/$Package/$Version"
$files = @($Filename, "$Filename.version", "$Filename.sha256")
foreach ($f in $files) {
Write-Host "[put] $api/$f"
curl.exe -sS -X PUT `
-H "Authorization: token $giteaToken" `
-H "Content-Type: application/octet-stream" `
"$api/$f" `
--data-binary "@$f"
if ($LASTEXITCODE -ne 0) { Write-Error "curl PUT $f завершился с кодом $LASTEXITCODE" }
}
# --- 6. Самопроверка: читаем .sha256 обратно и сверяем (контракт U4) ---
Write-Host "[verify] перечитываю $Filename.sha256 из Gitea..."
$remoteSum = (curl.exe -sS -H "Authorization: token $giteaToken" "$api/$Filename.sha256").Trim()
if ($remoteSum -ne $sha256Hex) {
Write-Error "Проверка не сошлась: remote=$remoteSum local=$sha256Hex"
}
$remoteVer = (curl.exe -sS -H "Authorization: token $giteaToken" "$api/$Filename.version").Trim()
if ($remoteVer -ne $Version) {
Write-Error "Версия не сошлась: remote=$remoteVer local=$Version"
}
Write-Host "[ok] опубликовано и проверено: $Version/$Filename (+.version, +.sha256)"