Compare commits
7 Commits
d9d043ec8e
...
feat/b9d90
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8cf4fc9f7c | ||
|
|
963e7b478e | ||
|
|
cd0619926e | ||
|
|
2854697415 | ||
|
|
ef812bb3d7 | ||
| 23bae9a68f | |||
|
|
70287140ec |
@@ -40,4 +40,15 @@
|
|||||||
## Версии
|
## Версии
|
||||||
- `app.Version` — семантическая major.minor.patch (ручной инкремент: patch=фиксы,
|
- `app.Version` — семантическая major.minor.patch (ручной инкремент: patch=фиксы,
|
||||||
minor=новая обратно-совместимая функциональность, major=несовместимые изменения). Сейчас 0.2.2.
|
minor=новая обратно-совместимая функциональность, major=несовместимые изменения). Сейчас 0.2.2.
|
||||||
- `main.version` (ldflag) — build-идентификатор `commit-<sha7>`, отдельно от app.Version.
|
- `main.version` (ldflag) — build-идентификатор `commit-<sha7>`, отдельно от app.Version.
|
||||||
|
|
||||||
|
## Форматирование (важно на Windows)
|
||||||
|
- Репо на Windows-чекауте: `core.autocrlf=true` → файлы в рабочей копии с CRLF; `gofmt -l`
|
||||||
|
на CRLF-копии шумит (глобально ~60 файлов). **Проверять формат только на LF-версии** (напр.
|
||||||
|
`git -c core.autocrlf=false clone` во временный каталог) — так дифы видны корректно.
|
||||||
|
- gofmt 1.26 форматирует doc-comments (`//` перед `go build ...` в `// go build` → пустая строка
|
||||||
|
`//`) и выравнивание структур; не все файлы отформатированы по новой версии (предсуществующе).
|
||||||
|
- CI (`go vet` + `go test`) формат не проверяет → gofmt-дифы не ломают сборку.
|
||||||
|
- **Known race в тест-харнессе app:** `e2eChannel.Send` (`internal/app/e2e_test.go:213`) пишет
|
||||||
|
`c.sent` без лока, тест читает с главной горутины → `-race` ловит в `TestE2ENotificationsOnTransitions`.
|
||||||
|
Путь `worker→Notify→Send`; предсуществует, ещё не чинили (fix — мьютекс в `e2eChannel`).
|
||||||
@@ -31,13 +31,15 @@ docs/ ui-spec.md — спека Fyne UI (слои, event-bus, fyne.Do
|
|||||||
## Ключевые инварианты
|
## Ключевые инварианты
|
||||||
|
|
||||||
- **App.New-сигнатура:** `App.New(configPath, version, updateToken string, noUI bool)` (4-й параметр — headless; cgo-вариант собирается только при наличии C-компилятора).
|
- **App.New-сигнатура:** `App.New(configPath, version, updateToken string, noUI bool)` (4-й параметр — headless; cgo-вариант собирается только при наличии C-компилятора).
|
||||||
|
- **chat.Router:** асинхронная обработка входящих — **ограниченный пул `chatWorkers=4` воркеров + FIFO-очереди per-user** (`userState`, `workerLoop`/`runUser`). Порядок сообщений одного UserID сохраняется (флаг `scheduled` → один активный воркер на пользователя); разные пользователи обрабатываются параллельно (до 4 одновременных LLM-вызовов). Backpressure по `jobs` блокирует только перегруженного пользователя, не весь long-poll. `Processed()`/`WaitProcessed()` — синхронизация тестов.
|
||||||
|
- **FreeChat** (app): `sessions map[uid]sessionID` защищён `sync.Mutex` (пишется из разных воркеров роутера).
|
||||||
- **UI:** окно — ещё одна реализация `chat.Channel` (присоединяется в Router). Core не трогает UI; обмен — событийная шина (events). Кнопка «Завершить» = полный выход (SetOnQuit→cancel→UI.Run возвращается); закрытие крестиком = сворачивание, Core живёт. UI собирается с `--noui`/без cgo.
|
- **UI:** окно — ещё одна реализация `chat.Channel` (присоединяется в Router). Core не трогает UI; обмен — событийная шина (events). Кнопка «Завершить» = полный выход (SetOnQuit→cancel→UI.Run возвращается); закрытие крестиком = сворачивание, Core живёт. UI собирается с `--noui`/без cgo.
|
||||||
- **Фазы аналитика (Decision.Phase):** `ask`, `propose`, `ready` (два последних обрабатываются одинаково в core), `abort`. Требования валидатора: ask — chat_reply/questions; propose — хотя бы одно изменённое поле; ready — без изменённых полей.
|
- **Фазы аналитика (Decision.Phase):** `ask`, `propose`, `ready` (два последних обрабатываются одинаково в core), `abort`. Требования валидатора: ask — chat_reply/questions; propose — хотя бы одно изменённое поле; ready — без изменённых полей.
|
||||||
- **Статусы задач (internal/model):** draft→collecting→ready→approved→running→success/failed/timeout + cancelled/aborted/closed (терминальные). UserID — chat.ID (одна активная задача на чат).
|
- **Статусы задач (internal/model):** draft→collecting→ready→approved→running→success/failed/timeout + cancelled/aborted/closed (терминальные). UserID — chat.ID (одна активная задача на чат).
|
||||||
- **Decider/Worker/Analyst/Reviewer:** Decider=analyst интерфейс; Worker — polling-планировщик; Reviewer проверяет diff dev-ветки (R1-R6), вердикт JSON {passed, critical_issues, solid_violations, comments}.
|
- **Decider/Worker/Analyst/Reviewer:** Decider=analyst интерфейс; Worker — polling-планировщик; Reviewer проверяет diff dev-ветки (R1-R6), вердикт JSON {passed, critical_issues, solid_violations, comments}.
|
||||||
- **gitops (worker):** worktree-режим; feature-ветка `feat/<taskTag>` от origin/main; push через http.extraHeader, токен Bearer.
|
- **gitops (worker):** worktree-режим; feature-ветка `feat/<taskTag>` от origin/main; push через http.extraHeader, токен Bearer.
|
||||||
- **Пути «всё рядом с .exe»:** db/worktree резолвятся от ExeDir; config.yaml — рядом с бинарём, фоллбэк cwd.
|
- **Пути «всё рядом с .exe»:** db/worktree резолвятся от ExeDir; config.yaml — рядом с бинарём, фоллбэк cwd.
|
||||||
- **Автообновление:** авто = только Check+уведомление; замена — по /update; версии в `commit-<sha7>/` (не `latest/`); Verify сверяет предprod-версию (binary+в.в) .
|
- **Автообновление:** авто = только Check+уведомление; замена — по /update; версии в `commit-<sha7>/` (не `latest/`); Verify сверяет контрольную сумму бинаря против `.sha256` той же версии. **Perf:** `ResolveLatest` проверяет наличие бинаря версии **HEAD-пробой без скачивания тела** (405/501 → fallback `GET Range: bytes=0-0`), сортировка версий по id убыв.; один общий `http.Client` (keep-alive). Ошибка U4 — только при несовпадении суммы (пустой/отсутствующий `.sha256` пропускает проверку — M1-известное замечание).
|
||||||
|
|
||||||
## opencode (v2 HTTP API, >= 1.18.18)
|
## opencode (v2 HTTP API, >= 1.18.18)
|
||||||
|
|
||||||
|
|||||||
14
README.md
14
README.md
@@ -131,7 +131,7 @@ update:
|
|||||||
|
|
||||||
## Интеграция с opencode (субагенты)
|
## Интеграция с opencode (субагенты)
|
||||||
|
|
||||||
Субагенты (analyst / dev / reviewer) запускаются через **headless** `opencode serve`
|
Субагенты (analyst / dev / reviewer / postmortem) запускаются через **headless** `opencode serve`
|
||||||
по **v2 HTTP API** (префикс `/api/*`). Требуемая версия opencode: **>= 1.18.18**
|
по **v2 HTTP API** (префикс `/api/*`). Требуемая версия opencode: **>= 1.18.18**
|
||||||
(сборки с v2 HTTP API). Старый бинарь, отвечающий только на `/global/health`,
|
(сборки с v2 HTTP API). Старый бинарь, отвечающий только на `/global/health`,
|
||||||
не подходит: healthcheck падает с понятной ошибкой (класс O1).
|
не подходит: healthcheck падает с понятной ошибкой (класс O1).
|
||||||
@@ -173,6 +173,18 @@ update:
|
|||||||
В `internal/core` фазы `propose` и `ready` обрабатываются одинаково (применить черновик,
|
В `internal/core` фазы `propose` и `ready` обрабатываются одинаково (применить черновик,
|
||||||
проверить репозитории, поставить `ready` и отдать резюме).
|
проверить репозитории, поставить `ready` и отдать резюме).
|
||||||
|
|
||||||
|
## Постмортем после failed/timeout
|
||||||
|
|
||||||
|
Когда задача завершилась `failed` или `timeout`, воркер дополнительно запускает
|
||||||
|
**постмортем-анализ** (`internal/worker/postmortem.go`, agent `postmortem`):
|
||||||
|
|
||||||
|
- анализирует сессии dev/reviewer (промпты и выводы из `traces`);
|
||||||
|
- оценивает законченность этапов и причины сбоя;
|
||||||
|
- сохраняет результат как trace `agent=postmortem` и шлёт владельцу уведомление
|
||||||
|
«🔍 анализ (после <статус>): почему так случилось / что сделать».
|
||||||
|
|
||||||
|
Статус задачи постмортем не меняет; сбои самого анализа не влияют на исход задачи.
|
||||||
|
|
||||||
## Автообновление из Gitea Packages
|
## Автообновление из Gitea Packages
|
||||||
|
|
||||||
Бинарь умеет сам себя обновлять из generic-пакета в Gitea. Модель:
|
Бинарь умеет сам себя обновлять из generic-пакета в Gitea. Модель:
|
||||||
|
|||||||
@@ -223,7 +223,12 @@
|
|||||||
отбрасываются (буфер ограничен);
|
отбрасываются (буфер ограничен);
|
||||||
- `Clear()` — очистить при недоступной задаче.
|
- `Clear()` — очистить при недоступной задаче.
|
||||||
- Композитор: поле ввода + кнопки команд, подключённые к `ui.Commands` через
|
- Композитор: поле ввода + кнопки команд, подключённые к `ui.Commands` через
|
||||||
`SetCommands(c)` (ввод → `SendText`, кнопки → Start/Approve/Skip/Cancel).
|
`SetCommands(c)` (ввод → `SendText`, кнопки → Start/Approve/Skip/Retry/Cancel).
|
||||||
|
- Кнопка «Перезапустить» — полный аналог Telegram-команды `/retry N`: отправляет
|
||||||
|
`/retry N` тем же путём, что и ручной ввод (Commands → Router → Core).
|
||||||
|
Активна только для вкладки с привязанной задачей; окно сообщает ID задачи
|
||||||
|
через `SetBoundTask(taskID)` (при создании вкладки в `addTab` и при привязке
|
||||||
|
свободной вкладки в `bindTaskSession`); `0` — кнопка неактивна.
|
||||||
- Доставка: окно (как chat.Channel) рендерит `Send`/`Ask`/историю через
|
- Доставка: окно (как chat.Channel) рендерит `Send`/`Ask`/историю через
|
||||||
`Append`; строка роли = pure `FormatRole(role)` (👤/🤖).
|
`Append`; строка роли = pure `FormatRole(role)` (👤/🤖).
|
||||||
- Обновление — на потоке Fyne; вызывающий уже внутри `fyne.Do`.
|
- Обновление — на потоке Fyne; вызывающий уже внутри `fyne.Do`.
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ var FS embed.FS
|
|||||||
//
|
//
|
||||||
// Держим в курсе: если добавить файл в каталог, он автоматически попадёт
|
// Держим в курсе: если добавить файл в каталог, он автоматически попадёт
|
||||||
// в FS, но для детерминированной распаковки список лучше дополнять здесь.
|
// в FS, но для детерминированной распаковки список лучше дополнять здесь.
|
||||||
var Names = []string{"analyst", "dev", "reviewer", "chat"}
|
var Names = []string{"analyst", "dev", "reviewer", "chat", "postmortem"}
|
||||||
|
|
||||||
// WriteTo распаковывает всех встроенных агентов в каталог dir/agents
|
// WriteTo распаковывает всех встроенных агентов в каталог dir/agents
|
||||||
// (создаёт его). Файлы перезаписываются — встроенная копия всегда актуальна.
|
// (создаёт его). Файлы перезаписываются — встроенная копия всегда актуальна.
|
||||||
|
|||||||
28
internal/agents/postmortem.md
Normal file
28
internal/agents/postmortem.md
Normal file
@@ -0,0 +1,28 @@
|
|||||||
|
---
|
||||||
|
name: postmortem
|
||||||
|
description: Постмортем-аналитик Ratatoskr — анализирует сессии dev/reviewer после failed/timeout задачи и даёт резюме: почему так и что сделать, чтобы не повторялось
|
||||||
|
mode: primary
|
||||||
|
---
|
||||||
|
|
||||||
|
Ты — постмортем-аналитик в конвейере Ratatoskr. Задача завершилась неудачей (failed) или таймаутом (timeout). Ты анализируешь, что пошло не так, и даёшь резюме, которое поможет не допускать этого впредь.
|
||||||
|
|
||||||
|
Тебе приходит промпт с:
|
||||||
|
- задачей (название, цель, критерии готовности AC, репозитории, итоговый статус);
|
||||||
|
- сессиями субагентов dev и/или reviewer: их статус (success/failed/timeout), промпт и вывод (output).
|
||||||
|
|
||||||
|
ПРАВИЛА:
|
||||||
|
1. Проанализируй сессии dev и reviewer: какие в них проблемы, насколько завершён каждый этап (разработка, ревью).
|
||||||
|
2. Оцени «законченность этапов»: что успел сделать dev, проверял ли reviewer весь дифф, были ли заблокированы работы.
|
||||||
|
3. Сделай вывод — **почему так случилось**: ошибка в задании, неясные AC, технический сбой, неорганизованная работа агента и т.п.
|
||||||
|
4. Дай рекомендации — «что сделать, чтобы этого не было»: как уточнять задачу, какие AC добавлять, какой контекст предавать агентам, какие этапы конвейера ужесточить.
|
||||||
|
5. ПИШИ СВОЙ ОТВЕТ **ПРОСТЫМ ТЕКСТОМ НА РУССКОМ ЯЗЫКЕ**, без JSON, без markdown-обёрток и лишней разметки.
|
||||||
|
|
||||||
|
Формат ответа (два обязательных блока, коротко и по делу):
|
||||||
|
|
||||||
|
Почему так случилось:
|
||||||
|
- <причина 1>
|
||||||
|
- <причина 2>
|
||||||
|
|
||||||
|
Что сделать, чтобы это не повторялось:
|
||||||
|
- <рекомендация 1>
|
||||||
|
- <рекомендация 2>
|
||||||
@@ -11,6 +11,7 @@ import (
|
|||||||
"os/signal"
|
"os/signal"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
|
"sync"
|
||||||
"syscall"
|
"syscall"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -52,7 +53,7 @@ const packageOwner = "kamelion"
|
|||||||
// не следует путать с build-идентификатором `main.version` (commit-<sha7>),
|
// не следует путать с build-идентификатором `main.version` (commit-<sha7>),
|
||||||
// который вшивается ldflag'ом и используется автообновлением. Здесь номер
|
// который вшивается ldflag'ом и используется автообновлением. Здесь номер
|
||||||
// поднимается вручную перед каждым релизом/публикацией новой сборки.
|
// поднимается вручную перед каждым релизом/публикацией новой сборки.
|
||||||
const Version = "0.2.2"
|
const Version = "0.3.0"
|
||||||
|
|
||||||
// App — собранный конвейер.
|
// App — собранный конвейер.
|
||||||
type App struct {
|
type App struct {
|
||||||
@@ -83,6 +84,7 @@ type App struct {
|
|||||||
type FreeChat struct {
|
type FreeChat struct {
|
||||||
Runner analyst.OpenCodeRunner
|
Runner analyst.OpenCodeRunner
|
||||||
Worktree string
|
Worktree string
|
||||||
|
mu sync.Mutex // защищает sessions (пишется из воркеров роутера)
|
||||||
sessions map[chat.UserID]string // uid → opencode sessionID
|
sessions map[chat.UserID]string // uid → opencode sessionID
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -95,12 +97,18 @@ func NewFreeChat(runner analyst.OpenCodeRunner, worktree string) *FreeChat {
|
|||||||
// и возвращает ответ модели. Сессия продолжается (resume по sessionID), поэтому
|
// и возвращает ответ модели. Сессия продолжается (resume по sessionID), поэтому
|
||||||
// каждая вкладка ведёт свой независимый диалог.
|
// каждая вкладка ведёт свой независимый диалог.
|
||||||
func (f *FreeChat) Chat(ctx context.Context, uid chat.UserID, text string) (string, error) {
|
func (f *FreeChat) Chat(ctx context.Context, uid chat.UserID, text string) (string, error) {
|
||||||
|
f.mu.Lock()
|
||||||
sid := f.sessions[uid]
|
sid := f.sessions[uid]
|
||||||
|
f.mu.Unlock()
|
||||||
|
|
||||||
res, err := f.Runner.Run(ctx, text, f.Worktree, "chat", sid)
|
res, err := f.Runner.Run(ctx, text, f.Worktree, "chat", sid)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
f.mu.Lock()
|
||||||
f.sessions[uid] = res.SessionID
|
f.sessions[uid] = res.SessionID
|
||||||
|
f.mu.Unlock()
|
||||||
return res.Stdout, nil
|
return res.Stdout, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -21,36 +21,90 @@ type Router struct {
|
|||||||
// Hook, вызываемый на каждое входящее событие (обычно → process_turn).
|
// Hook, вызываемый на каждое входящее событие (обычно → process_turn).
|
||||||
onUserMsg func(Incoming)
|
onUserMsg func(Incoming)
|
||||||
|
|
||||||
// Асинхронная обработка входящих: handleIncoming кладёт событие в канал,
|
// Асинхронная обработка входящих ограниченным пулом воркеров с
|
||||||
// воркер-горутина последовательно вызывает onUserMsg. Благодаря этому
|
// упорядоченными очередями per-user (см. chatWorkers, userState).
|
||||||
// long-poll цикл канала (Telegram) не блокируется на время долгого
|
// Благодаря этому long-poll цикл канала (Telegram) не блокируется на время
|
||||||
// вызова аналитика и продолжает принимать новые сообщения.
|
// долгого вызова аналитика, а сообщения разных пользователей не сериализуются
|
||||||
incoming chan Incoming
|
// друг за другом: каждый активный пользователь занимает своего воркера.
|
||||||
|
jobs chan *userState
|
||||||
|
users map[UserID]*userState
|
||||||
|
userMu sync.Mutex
|
||||||
|
|
||||||
// processed — число обработанных воркером событий (для синхронизации
|
// processed — число обработанных воркером событий (для синхронизации
|
||||||
// тестов с асинхронной очередью: WaitProcessed ждёт обработку события).
|
// тестов с асинхронной очередью: WaitProcessed ждёт обработку события).
|
||||||
processed atomic.Int64
|
processed atomic.Int64
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// userState — FIFO-очередь входящих одного пользователя. В каждый момент
|
||||||
|
// для пользователя активен ровно один воркер (scheduled), поэтому порядок
|
||||||
|
// обработки его сообщений сохраняется, а параллелизм достигается между
|
||||||
|
// разными пользователями.
|
||||||
|
type userState struct {
|
||||||
|
mu sync.Mutex
|
||||||
|
pending []Incoming
|
||||||
|
scheduled bool
|
||||||
|
}
|
||||||
|
|
||||||
|
// chatWorkers — число воркеров обработки входящих. Ограничивает количество
|
||||||
|
// одновременных тяжёлых LLM-вызовов (аналитик/свободный чат), чтобы поток
|
||||||
|
// каналов не упирался в один долгий вызов.
|
||||||
|
const chatWorkers = 4
|
||||||
|
|
||||||
// NewRouter создаёт роутер. onUserMsg — колбэк обработки входящего.
|
// NewRouter создаёт роутер. onUserMsg — колбэк обработки входящего.
|
||||||
func NewRouter(onUserMsg func(Incoming)) *Router {
|
func NewRouter(onUserMsg func(Incoming)) *Router {
|
||||||
if onUserMsg == nil {
|
if onUserMsg == nil {
|
||||||
onUserMsg = func(Incoming) {}
|
onUserMsg = func(Incoming) {}
|
||||||
}
|
}
|
||||||
r := &Router{
|
r := &Router{
|
||||||
sessions: map[UserID]any{},
|
sessions: map[UserID]any{},
|
||||||
routes: map[UserID]Route{},
|
routes: map[UserID]Route{},
|
||||||
pending: map[UserID]PendingQ{},
|
pending: map[UserID]PendingQ{},
|
||||||
onUserMsg: onUserMsg,
|
onUserMsg: onUserMsg,
|
||||||
incoming: make(chan Incoming, 256),
|
jobs: make(chan *userState, chatWorkers),
|
||||||
|
users: make(map[UserID]*userState),
|
||||||
|
}
|
||||||
|
for i := 0; i < chatWorkers; i++ {
|
||||||
|
go r.workerLoop()
|
||||||
}
|
}
|
||||||
go r.processLoop()
|
|
||||||
return r
|
return r
|
||||||
}
|
}
|
||||||
|
|
||||||
// processLoop — воркер асинхронной обработки входящих (FIFO).
|
// userStateOf возвращает очередь пользователя (создаёт при первом сообщении).
|
||||||
func (r *Router) processLoop() {
|
// Очереди живут вечно — по одной маленькой структуре на пользователя/вкладку.
|
||||||
for inc := range r.incoming {
|
func (r *Router) userStateOf(uid UserID) *userState {
|
||||||
|
r.userMu.Lock()
|
||||||
|
defer r.userMu.Unlock()
|
||||||
|
st, ok := r.users[uid]
|
||||||
|
if !ok {
|
||||||
|
st = &userState{}
|
||||||
|
r.users[uid] = st
|
||||||
|
}
|
||||||
|
return st
|
||||||
|
}
|
||||||
|
|
||||||
|
// workerLoop — воркер пула: берёт пользователя из общей очереди и дренит его.
|
||||||
|
func (r *Router) workerLoop() {
|
||||||
|
for st := range r.jobs {
|
||||||
|
r.runUser(st)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// runUser обрабатывает все накопленные сообщения пользователя по порядку.
|
||||||
|
// По исчерпании очереди снимает scheduled — следующий handleIncoming вновь
|
||||||
|
// поставит пользователя в jobs. Возвращается в workerLoop, чтобы тот взял
|
||||||
|
// следующего пользователя из общей очереди.
|
||||||
|
func (r *Router) runUser(st *userState) {
|
||||||
|
for {
|
||||||
|
st.mu.Lock()
|
||||||
|
if len(st.pending) == 0 {
|
||||||
|
st.scheduled = false
|
||||||
|
st.mu.Unlock()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
inc := st.pending[0]
|
||||||
|
st.pending = st.pending[1:]
|
||||||
|
st.mu.Unlock()
|
||||||
|
|
||||||
r.onUserMsg(inc)
|
r.onUserMsg(inc)
|
||||||
r.processed.Add(1)
|
r.processed.Add(1)
|
||||||
}
|
}
|
||||||
@@ -99,9 +153,19 @@ func (r *Router) handleIncoming(inc Incoming) {
|
|||||||
}
|
}
|
||||||
r.mu.Unlock()
|
r.mu.Unlock()
|
||||||
|
|
||||||
// Асинхронная обработка: кладём событие в очередь воркера и сразу
|
// Асинхронная обработка: кладём событие в FIFO-очередь пользователя и
|
||||||
// возвращаемся, не блокируя вызывающий long-poll цикл канала.
|
// сразу возвращаемся, не блокируя вызывающий long-poll цикл канала.
|
||||||
r.incoming <- inc
|
// Если пользователь ещё не обрабатывается — ставим его в общую очередь
|
||||||
|
// пула воркеров. Backpressure по jobs блокирует только перегруженного
|
||||||
|
// пользователя (его собственную горутину канала), не весь роутер.
|
||||||
|
st := r.userStateOf(inc.UserID)
|
||||||
|
st.mu.Lock()
|
||||||
|
st.pending = append(st.pending, inc)
|
||||||
|
if !st.scheduled {
|
||||||
|
st.scheduled = true
|
||||||
|
r.jobs <- st
|
||||||
|
}
|
||||||
|
st.mu.Unlock()
|
||||||
}
|
}
|
||||||
|
|
||||||
// Send уведомляет пользователя через текущий маршрут. M1 (нет маршрута) — no-op,
|
// Send уведомляет пользователя через текущий маршрут. M1 (нет маршрута) — no-op,
|
||||||
|
|||||||
@@ -9,8 +9,8 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
uidA UserID = "u-a"
|
uidA UserID = "u-a"
|
||||||
uidB UserID = "u-b"
|
uidB UserID = "u-b"
|
||||||
tg Address = "tg://123"
|
tg Address = "tg://123"
|
||||||
tui Address = "tui://local"
|
tui Address = "tui://local"
|
||||||
)
|
)
|
||||||
@@ -213,3 +213,105 @@ func TestRouter_Ask_PendingNotConsumedFromOtherAddr(t *testing.T) {
|
|||||||
t.Fatal("pending должен остаться (ответ из другого адреса)")
|
t.Fatal("pending должен остаться (ответ из другого адреса)")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestRouter_PerUserOrdering проверяет, что сообщения одного пользователя
|
||||||
|
// обрабатываются строго в порядке поступления (пул воркеров не перемешивает).
|
||||||
|
func TestRouter_PerUserOrdering(t *testing.T) {
|
||||||
|
cb := newFakeOnMsg()
|
||||||
|
r := NewRouter(cb.h)
|
||||||
|
tgCh := newFakeChannel(tg)
|
||||||
|
_ = r.Attach(tgCh)
|
||||||
|
|
||||||
|
for _, txt := range []string{"1", "2", "3"} {
|
||||||
|
tgCh.emit(uidA, tg, txt)
|
||||||
|
}
|
||||||
|
if !cb.wait(3) {
|
||||||
|
t.Fatal("сообщения не обработаны за таймаут")
|
||||||
|
}
|
||||||
|
for i, want := range []string{"1", "2", "3"} {
|
||||||
|
if got := cb.get(i).Msg.Text; got != want {
|
||||||
|
t.Errorf("порядок обработки нарушен: idx %d = %q, want %q", i, got, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRouter_ParallelismAcrossUsers проверяет, что пока обработчик одного
|
||||||
|
// пользователя заблокирован (долгий LLM-вызов), сообщение другого пользователя
|
||||||
|
// обрабатывается в другом воркере, а второе сообщение того же пользователя —
|
||||||
|
// ждёт своей очереди (per-user порядок).
|
||||||
|
func TestRouter_ParallelismAcrossUsers(t *testing.T) {
|
||||||
|
r := NewRouter(nil)
|
||||||
|
|
||||||
|
tgCh := newFakeChannel(tg)
|
||||||
|
tuiCh := newFakeChannel(tui)
|
||||||
|
_ = r.Attach(tgCh)
|
||||||
|
_ = r.Attach(tuiCh)
|
||||||
|
|
||||||
|
blocked := make(chan struct{})
|
||||||
|
release := make(chan struct{})
|
||||||
|
var muLocal sync.Mutex
|
||||||
|
seen := make([]string, 0, 3)
|
||||||
|
signal := make(chan struct{}, 8)
|
||||||
|
|
||||||
|
h := func(inc Incoming) {
|
||||||
|
if inc.Msg.Text == "block" {
|
||||||
|
close(blocked)
|
||||||
|
<-release // держим воркера, пока не отпустим
|
||||||
|
}
|
||||||
|
muLocal.Lock()
|
||||||
|
seen = append(seen, inc.Msg.Text)
|
||||||
|
muLocal.Unlock()
|
||||||
|
signal <- struct{}{}
|
||||||
|
}
|
||||||
|
r.onUserMsg = h
|
||||||
|
|
||||||
|
snapshot := func() []string {
|
||||||
|
muLocal.Lock()
|
||||||
|
defer muLocal.Unlock()
|
||||||
|
return append([]string(nil), seen...)
|
||||||
|
}
|
||||||
|
waitFor := func(n int) bool {
|
||||||
|
deadline := time.After(2 * time.Second)
|
||||||
|
for len(snapshot()) < n {
|
||||||
|
select {
|
||||||
|
case <-signal:
|
||||||
|
case <-deadline:
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
// первое сообщение A блокирует своего воркера
|
||||||
|
tgCh.emit(uidA, tg, "block")
|
||||||
|
<-blocked
|
||||||
|
|
||||||
|
// B обрабатывается параллельно, пока A висит
|
||||||
|
tuiCh.emit(uidB, tui, "B1")
|
||||||
|
select {
|
||||||
|
case <-signal:
|
||||||
|
case <-time.After(100 * time.Millisecond):
|
||||||
|
t.Fatal("B не обработан, пока блокирован A — чат-путь снова сериализован")
|
||||||
|
}
|
||||||
|
if got := snapshot(); len(got) != 1 || got[0] != "B1" {
|
||||||
|
t.Fatalf("ожидали обработку B1, got %v", got)
|
||||||
|
}
|
||||||
|
|
||||||
|
// второе сообщение A НЕ обрабатывается, пока занят воркер A (порядок per-user)
|
||||||
|
tgCh.emit(uidA, tg, "a2")
|
||||||
|
select {
|
||||||
|
case <-signal:
|
||||||
|
t.Fatal("сообщение A обработано ДО освобождения A — нарушен per-user порядок")
|
||||||
|
case <-time.After(80 * time.Millisecond):
|
||||||
|
}
|
||||||
|
|
||||||
|
// отпускаем A → дообрабатывается a2
|
||||||
|
close(release)
|
||||||
|
if !waitFor(3) {
|
||||||
|
t.Fatal("итоговые сообщения не обработаны")
|
||||||
|
}
|
||||||
|
got := snapshot()
|
||||||
|
if got[2] != "a2" {
|
||||||
|
t.Errorf("порядок персональной очереди нарушен: pos2 = %q, want a2", got[2])
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -85,6 +85,8 @@ func TestClassifyLevel(t *testing.T) {
|
|||||||
{"app: db opened /tmp/r.db", LevelInfo},
|
{"app: db opened /tmp/r.db", LevelInfo},
|
||||||
{"app: worker started", LevelInfo},
|
{"app: worker started", LevelInfo},
|
||||||
{"opencode: debug: poll request", LevelDebug},
|
{"opencode: debug: poll request", LevelDebug},
|
||||||
|
{"opencode api debug: prompt -> POST http://127.0.0.1:4096/api/session", LevelDebug},
|
||||||
|
{"opencode api debug: messages response (512 bytes)", LevelDebug},
|
||||||
{"trace: session resumed", LevelDebug},
|
{"trace: session resumed", LevelDebug},
|
||||||
{"tg: warn: long poll timeout", LevelWarning},
|
{"tg: warn: long poll timeout", LevelWarning},
|
||||||
{"ПРЕДУПРЕЖДЕНИЕ: конфиг не задан", LevelWarning},
|
{"ПРЕДУПРЕЖДЕНИЕ: конфиг не задан", LevelWarning},
|
||||||
@@ -119,4 +121,4 @@ func TestLogWriterClassifiesLevels(t *testing.T) {
|
|||||||
if second.Level != LevelError || second.Text != "ERROR: boom" {
|
if second.Level != LevelError || second.Text != "ERROR: boom" {
|
||||||
t.Fatalf("unexpected second: %+v", second)
|
t.Fatalf("unexpected second: %+v", second)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -24,9 +24,9 @@ import (
|
|||||||
// Prompt не блокирует: вердикт собирается поллингом из content[].type=="text"
|
// Prompt не блокирует: вердикт собирается поллингом из content[].type=="text"
|
||||||
// новых assistant-сообщений (см. Runner.awaitVerdict).
|
// новых 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 // единый клиент: все операции быстрые (нет блокирующего Send)
|
http *http.Client // единый клиент: все операции быстрые (нет блокирующего Send)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -63,10 +63,7 @@ func (c *Client) do(ctx context.Context, method, path, op string, body []byte) (
|
|||||||
req.Header.Set("Content-Type", "application/json")
|
req.Header.Set("Content-Type", "application/json")
|
||||||
}
|
}
|
||||||
if c.Debug {
|
if c.Debug {
|
||||||
log.Printf("opencode api %s -> %s %s%s", op, method, c.BaseURL, path)
|
log.Printf("opencode api debug: %s -> %s %s%s", op, method, c.BaseURL, path)
|
||||||
if len(body) > 0 {
|
|
||||||
log.Printf("opencode api %s request body: %s", op, truncateStr(string(body), 5000))
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
resp, err := c.http.Do(req)
|
resp, err := c.http.Do(req)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -79,12 +76,12 @@ func (c *Client) do(ctx context.Context, method, path, op string, body []byte) (
|
|||||||
}
|
}
|
||||||
if resp.StatusCode < 200 || resp.StatusCode > 299 {
|
if resp.StatusCode < 200 || resp.StatusCode > 299 {
|
||||||
if c.Debug {
|
if c.Debug {
|
||||||
log.Printf("opencode api %s response: status %d: %s", op, resp.StatusCode, truncateStr(string(b), 1000))
|
log.Printf("opencode api debug: %s response: status %d", op, resp.StatusCode)
|
||||||
}
|
}
|
||||||
return nil, &ClientErr{Op: op, Err: fmt.Errorf("status %d: %s", resp.StatusCode, truncateStr(string(b), 300))}
|
return nil, &ClientErr{Op: op, Err: fmt.Errorf("status %d: %s", resp.StatusCode, truncateStr(string(b), 300))}
|
||||||
}
|
}
|
||||||
if c.Debug {
|
if c.Debug {
|
||||||
log.Printf("opencode api %s response (%d bytes): %s", op, len(b), truncateStr(string(b), 5000))
|
log.Printf("opencode api debug: %s response (%d bytes)", op, len(b))
|
||||||
}
|
}
|
||||||
return b, nil
|
return b, nil
|
||||||
}
|
}
|
||||||
@@ -323,6 +320,26 @@ func assistantText(msgs []v2Message, since int64) []string {
|
|||||||
return texts
|
return texts
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// assistantVerdict собирает финальный текст ответа: сначала text-парты, а если
|
||||||
|
// их нет — только reasoning-парты (fallback для моделей, которые на некоторые
|
||||||
|
// запросы отвечают лишь reasoning без text). usedReasoning=true означает, что
|
||||||
|
// text-партов не было вовсе и вердикт собран из reasoning.
|
||||||
|
func assistantVerdict(msgs []v2Message, since int64) (texts []string, usedReasoning bool) {
|
||||||
|
if texts := assistantText(msgs, since); len(texts) > 0 {
|
||||||
|
return texts, false
|
||||||
|
}
|
||||||
|
ass := assistantSince(msgs, since)
|
||||||
|
reasoning := make([]string, 0, len(ass))
|
||||||
|
for i := len(ass) - 1; i >= 0; i-- {
|
||||||
|
for _, p := range ass[i].Content {
|
||||||
|
if p.Type == "reasoning" && p.Text != "" {
|
||||||
|
reasoning = append(reasoning, p.Text)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return reasoning, len(reasoning) > 0
|
||||||
|
}
|
||||||
|
|
||||||
func truncateStr(s string, n int) string {
|
func truncateStr(s string, n int) string {
|
||||||
if len(s) <= n {
|
if len(s) <= n {
|
||||||
return s
|
return s
|
||||||
|
|||||||
@@ -25,10 +25,11 @@ type fakeAPIServer struct {
|
|||||||
sessionID string
|
sessionID string
|
||||||
created bool
|
created bool
|
||||||
active bool
|
active bool
|
||||||
blockPrompt bool
|
blockPrompt bool
|
||||||
messages []v2Message
|
messages []v2Message
|
||||||
verdictText string
|
verdictText string
|
||||||
failCreate bool
|
verdictReasoning string // завершённый ответ только с reasoning-партом (без text)
|
||||||
|
failCreate bool
|
||||||
failMessages bool
|
failMessages bool
|
||||||
createdModel *ModelRef // модель, полученная на POST /api/session
|
createdModel *ModelRef // модель, полученная на POST /api/session
|
||||||
promptCalls int
|
promptCalls int
|
||||||
@@ -133,6 +134,9 @@ func (f *fakeAPIServer) handler() http.Handler {
|
|||||||
if msgs == nil && f.verdictText != "" && !f.blockPrompt {
|
if msgs == nil && f.verdictText != "" && !f.blockPrompt {
|
||||||
msgs = []v2Message{f.assistantMsg(f.verdictText)}
|
msgs = []v2Message{f.assistantMsg(f.verdictText)}
|
||||||
}
|
}
|
||||||
|
if msgs == nil && f.verdictReasoning != "" && !f.blockPrompt {
|
||||||
|
msgs = []v2Message{f.assistantReasoningMsg(f.verdictReasoning)}
|
||||||
|
}
|
||||||
if msgs == nil {
|
if msgs == nil {
|
||||||
msgs = []v2Message{}
|
msgs = []v2Message{}
|
||||||
}
|
}
|
||||||
@@ -153,6 +157,19 @@ func (f *fakeAPIServer) assistantMsg(text string) v2Message {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// assistantReasoningMsg строит завершённое assistant-сообщение только с
|
||||||
|
// reasoning-партом (без text) — для проверки fallback-сценария.
|
||||||
|
func (f *fakeAPIServer) assistantReasoningMsg(text string) v2Message {
|
||||||
|
now := time.Now().UnixMilli()
|
||||||
|
return v2Message{
|
||||||
|
ID: "msg_r",
|
||||||
|
Type: "assistant",
|
||||||
|
Content: []v2Part{{Type: "reasoning", 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)
|
||||||
@@ -279,6 +296,48 @@ func Test_newestAssistant(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func Test_assistantVerdict(t *testing.T) {
|
||||||
|
older := time.Now().Add(-time.Minute).UnixMilli()
|
||||||
|
newer := time.Now().UnixMilli()
|
||||||
|
reasoningOf := func(text string, at *int64) v2Message {
|
||||||
|
return v2Message{ID: "r", Type: "assistant", Content: []v2Part{{Type: "reasoning", Text: text}}, Time: v2Time{Created: at}}
|
||||||
|
}
|
||||||
|
|
||||||
|
// reasoning-only: text-партов нет → fallback на reasoning, usedReasoning=true.
|
||||||
|
// Сообщения приходят новейшими первыми (как из API) → размышление 2 новее.
|
||||||
|
reasoningOnlyMsgs := []v2Message{
|
||||||
|
reasoningOf("размышление 2", &newer),
|
||||||
|
reasoningOf("размышление 1", &older),
|
||||||
|
}
|
||||||
|
texts, used := assistantVerdict(reasoningOnlyMsgs, older)
|
||||||
|
if !used {
|
||||||
|
t.Error("usedReasoning = false, want true для reasoning-only")
|
||||||
|
}
|
||||||
|
if len(texts) != 2 || texts[0] != "размышление 1" || texts[1] != "размышление 2" {
|
||||||
|
t.Errorf("verdict = %v, want [размышление 1 размышление 2] (хронологически)", texts)
|
||||||
|
}
|
||||||
|
|
||||||
|
// text + reasoning → берётся text, reasoning игнорируется.
|
||||||
|
mixed := []v2Message{
|
||||||
|
{ID: "a", Type: "assistant",
|
||||||
|
Content: []v2Part{{Type: "reasoning", Text: "thinking"}, {Type: "text", Text: "ответ"}},
|
||||||
|
Time: v2Time{Created: &newer}},
|
||||||
|
}
|
||||||
|
texts, used = assistantVerdict(mixed, older)
|
||||||
|
if used {
|
||||||
|
t.Error("usedReasoning = true, want false (есть text)")
|
||||||
|
}
|
||||||
|
if len(texts) != 1 || texts[0] != "ответ" {
|
||||||
|
t.Errorf("verdict = %v, want [ответ]", texts)
|
||||||
|
}
|
||||||
|
|
||||||
|
// пусто → пусто и usedReasoning=false.
|
||||||
|
empty := []v2Message{{ID: "u", Type: "user", Time: v2Time{Created: &newer}}}
|
||||||
|
if texts, used := assistantVerdict(empty, older); used || len(texts) != 0 {
|
||||||
|
t.Errorf("пусто: texts=%v usedReasoning=%v, want пусто/false", texts, used)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func Test_parseModelString(t *testing.T) {
|
func Test_parseModelString(t *testing.T) {
|
||||||
m := parseModelString("tokentool/deepseek/deepseek-v4-flash-0731")
|
m := parseModelString("tokentool/deepseek/deepseek-v4-flash-0731")
|
||||||
if m == nil || m.ProviderID != "tokentool" || m.ID != "deepseek/deepseek-v4-flash-0731" {
|
if m == nil || m.ProviderID != "tokentool" || m.ID != "deepseek/deepseek-v4-flash-0731" {
|
||||||
|
|||||||
@@ -216,10 +216,13 @@ func (r *Runner) verdict(model *ModelRef, cur *v2Message, msgs []v2Message, sinc
|
|||||||
if cur.Error != nil && cur.Error.Message != "" {
|
if cur.Error != nil && cur.Error.Message != "" {
|
||||||
return nil, &ClientErr{Op: "prompt", Err: errors.New(cur.Error.Message)}
|
return nil, &ClientErr{Op: "prompt", Err: errors.New(cur.Error.Message)}
|
||||||
}
|
}
|
||||||
texts := assistantText(msgs, since)
|
texts, usedReasoning := assistantVerdict(msgs, since)
|
||||||
if len(texts) == 0 {
|
if len(texts) == 0 {
|
||||||
return nil, &ClientErr{Op: "prompt", Err: errors.New("нет text-части в ответе")}
|
return nil, &ClientErr{Op: "prompt", Err: errors.New("нет text-части в ответе")}
|
||||||
}
|
}
|
||||||
|
if usedReasoning {
|
||||||
|
r.logf("WARN opencode: в ответе нет text-части — использую reasoning-парты как вердикт")
|
||||||
|
}
|
||||||
vd := stripFence(strings.Join(texts, "\n"))
|
vd := stripFence(strings.Join(texts, "\n"))
|
||||||
r.logf("opencode вердикт готов (%d байт)", len(vd))
|
r.logf("opencode вердикт готов (%d байт)", len(vd))
|
||||||
return &Result{RC: 0, Stdout: vd, SessionID: sid}, nil
|
return &Result{RC: 0, Stdout: vd, SessionID: sid}, nil
|
||||||
|
|||||||
@@ -110,6 +110,28 @@ func TestRun_ReasoningGrowth(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestRun_ReasoningOnlyVerdict(t *testing.T) {
|
||||||
|
t.Setenv("XDG_CONFIG_HOME", t.TempDir())
|
||||||
|
dir := t.TempDir()
|
||||||
|
// завершённый ответ без text-парта, только reasoning — вердикт собирается
|
||||||
|
// из reasoning (fallback) вместо ошибки «нет text-части в ответе».
|
||||||
|
f := &fakeAPIServer{verdictReasoning: "размышления без текста"}
|
||||||
|
p, _ := fakePool(t, f, dir)
|
||||||
|
|
||||||
|
r := &Runner{Pool: p, IdleTimeout: time.Minute, HardTimeout: time.Minute,
|
||||||
|
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-only вердикт)", res.RC)
|
||||||
|
}
|
||||||
|
if !contains(res.Stdout, "размышления без текста") {
|
||||||
|
t.Errorf("Stdout = %q, want reasoning fallback", res.Stdout)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestRun_ContextCancel(t *testing.T) {
|
func TestRun_ContextCancel(t *testing.T) {
|
||||||
t.Setenv("XDG_CONFIG_HOME", t.TempDir())
|
t.Setenv("XDG_CONFIG_HOME", t.TempDir())
|
||||||
dir := t.TempDir()
|
dir := t.TempDir()
|
||||||
|
|||||||
@@ -15,6 +15,10 @@ type ChatPanel interface {
|
|||||||
Clear()
|
Clear()
|
||||||
// SetCommands подключает команды к композитору (ввод + кнопки).
|
// SetCommands подключает команды к композитору (ввод + кнопки).
|
||||||
SetCommands(c *Commands)
|
SetCommands(c *Commands)
|
||||||
|
// SetBoundTask сообщает панели задачу, к которой привязана вкладка
|
||||||
|
// (0 — свободный чат). Кнопка «Перезапустить» (/retry N) активна только
|
||||||
|
// для вкладки с привязанной задачей.
|
||||||
|
SetBoundTask(taskID int64)
|
||||||
}
|
}
|
||||||
|
|
||||||
// NilChatPanel — no-op реализация ChatPanel для headless-режима и тестов.
|
// NilChatPanel — no-op реализация ChatPanel для headless-режима и тестов.
|
||||||
@@ -27,6 +31,7 @@ func (NilChatPanel) SetTranscript(string) {}
|
|||||||
func (NilChatPanel) Append(string) {}
|
func (NilChatPanel) Append(string) {}
|
||||||
func (NilChatPanel) Clear() {}
|
func (NilChatPanel) Clear() {}
|
||||||
func (NilChatPanel) SetCommands(*Commands) {}
|
func (NilChatPanel) SetCommands(*Commands) {}
|
||||||
|
func (NilChatPanel) SetBoundTask(int64) {}
|
||||||
|
|
||||||
// FormatRole — строка роли в диалоге (pure-функция, спец 12.12).
|
// FormatRole — строка роли в диалоге (pure-функция, спец 12.12).
|
||||||
func FormatRole(role string) string {
|
func FormatRole(role string) string {
|
||||||
|
|||||||
@@ -23,4 +23,32 @@ func TestNilChatPanelNoop(t *testing.T) {
|
|||||||
p.Append("строка")
|
p.Append("строка")
|
||||||
p.Clear()
|
p.Clear()
|
||||||
p.SetCommands(c)
|
p.SetCommands(c)
|
||||||
|
p.SetBoundTask(7)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Кнопка «Перезапустить» — полный аналог Telegram-команды /retry N:
|
||||||
|
// в канал уходит ровно тот же текст, что при ручном вводе команды.
|
||||||
|
func TestRetryButtonSubmitsSameCommandAsTelegram(t *testing.T) {
|
||||||
|
var got []string
|
||||||
|
c := NewCommands(func(text string) { got = append(got, text) })
|
||||||
|
|
||||||
|
c.Retry(42)
|
||||||
|
|
||||||
|
if len(got) != 1 || got[0] != "/retry 42" {
|
||||||
|
t.Fatalf("got %v, want [/retry 42]", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Свободная вкладка (без задачи): Retry(id<=0) — no-op, ничего не отправляется
|
||||||
|
// (кнопка «Перезапустить» на такой вкладке неактивна).
|
||||||
|
func TestRetryWithoutTaskIsNoop(t *testing.T) {
|
||||||
|
var got []string
|
||||||
|
c := NewCommands(func(text string) { got = append(got, text) })
|
||||||
|
|
||||||
|
c.Retry(0)
|
||||||
|
c.Retry(-5)
|
||||||
|
|
||||||
|
if len(got) != 0 {
|
||||||
|
t.Fatalf("got %v, want nothing submitted", got)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
@@ -30,8 +30,14 @@ func (c *Commands) Cancel() { c.Submit("/cancel") }
|
|||||||
// Skip — пропустить сбор, сформировать черновик (/skip).
|
// Skip — пропустить сбор, сформировать черновик (/skip).
|
||||||
func (c *Commands) Skip() { c.Submit("/skip") }
|
func (c *Commands) Skip() { c.Submit("/skip") }
|
||||||
|
|
||||||
// Retry — перезапустить задачу N (/retry N).
|
// Retry — перезапустить задачу N (/retry N): полный аналог Telegram-команды.
|
||||||
func (c *Commands) Retry(id int64) { c.Submit(fmt.Sprintf("/retry %d", id)) }
|
// id <= 0 — no-op (кнопка «Перезапустить» на свободной вкладке неактивна).
|
||||||
|
func (c *Commands) Retry(id int64) {
|
||||||
|
if id <= 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
c.Submit(fmt.Sprintf("/retry %d", id))
|
||||||
|
}
|
||||||
|
|
||||||
// Status — запросить статус задачи N (/status N).
|
// Status — запросить статус задачи N (/status N).
|
||||||
func (c *Commands) Status(id int64) { c.Submit(fmt.Sprintf("/status %d", id)) }
|
func (c *Commands) Status(id int64) { c.Submit(fmt.Sprintf("/status %d", id)) }
|
||||||
|
|||||||
@@ -18,9 +18,11 @@ const chatLimit = 200_000
|
|||||||
// ChatPanel — Fyne-реализация ui.ChatPanel (спец 12.8, 12.12): транскрипт
|
// ChatPanel — Fyne-реализация ui.ChatPanel (спец 12.8, 12.12): транскрипт
|
||||||
// диалога + композитор (поле ввода и кнопки команд).
|
// диалога + композитор (поле ввода и кнопки команд).
|
||||||
type ChatPanel struct {
|
type ChatPanel struct {
|
||||||
convLbl *widget.Label
|
convLbl *widget.Label
|
||||||
input *widget.Entry
|
input *widget.Entry
|
||||||
commands *ui.Commands
|
commands *ui.Commands
|
||||||
|
retryBtn *widget.Button // «Перезапустить» — полный аналог /retry N
|
||||||
|
boundTask int64 // >0 — вкладка привязана к задаче №boundTask
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewChatPanel создаёт ChatPanel.
|
// NewChatPanel создаёт ChatPanel.
|
||||||
@@ -30,6 +32,8 @@ func NewChatPanel() *ChatPanel {
|
|||||||
p.convLbl.Wrapping = fyne.TextWrapWord
|
p.convLbl.Wrapping = fyne.TextWrapWord
|
||||||
p.input = widget.NewEntry()
|
p.input = widget.NewEntry()
|
||||||
p.input.SetPlaceHolder("Сообщение… (Enter — отправить)")
|
p.input.SetPlaceHolder("Сообщение… (Enter — отправить)")
|
||||||
|
p.retryBtn = widget.NewButton("Перезапустить", p.retryTask)
|
||||||
|
p.retryBtn.Disable() // активируется при привязке вкладки к задаче (SetBoundTask)
|
||||||
return p
|
return p
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -44,11 +48,32 @@ func (p *ChatPanel) Composer() fyne.CanvasObject {
|
|||||||
widget.NewButton("Новая", func() { if p.commands != nil { p.commands.Start() } }),
|
widget.NewButton("Новая", func() { if p.commands != nil { p.commands.Start() } }),
|
||||||
widget.NewButton("Создавай", func() { if p.commands != nil { p.commands.Approve() } }),
|
widget.NewButton("Создавай", func() { if p.commands != nil { p.commands.Approve() } }),
|
||||||
widget.NewButton("Пропустить", func() { if p.commands != nil { p.commands.Skip() } }),
|
widget.NewButton("Пропустить", func() { if p.commands != nil { p.commands.Skip() } }),
|
||||||
|
p.retryBtn,
|
||||||
widget.NewButton("Отмена", func() { if p.commands != nil { p.commands.Cancel() } }),
|
widget.NewButton("Отмена", func() { if p.commands != nil { p.commands.Cancel() } }),
|
||||||
)
|
)
|
||||||
return container.NewBorder(nil, nil, cmdBar, nil, p.input)
|
return container.NewBorder(nil, nil, cmdBar, nil, p.input)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// retryTask — колбэк кнопки «Перезапустить»: отправляет "/retry N" тем же
|
||||||
|
// путём, что и ввод пользователя (Commands → Router → Core), т.е. серверная
|
||||||
|
// логика полностью совпадает с Telegram-командой /retry N.
|
||||||
|
func (p *ChatPanel) retryTask() {
|
||||||
|
if p.commands != nil {
|
||||||
|
p.commands.Retry(p.boundTask)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// SetBoundTask включает кнопку «Перезапустить», когда вкладка привязана к
|
||||||
|
// задаче (taskID > 0); для свободной вкладки кнопка неактивна.
|
||||||
|
func (p *ChatPanel) SetBoundTask(taskID int64) {
|
||||||
|
p.boundTask = taskID
|
||||||
|
if taskID > 0 {
|
||||||
|
p.retryBtn.Enable()
|
||||||
|
} else {
|
||||||
|
p.retryBtn.Disable()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// SetTranscript заменяет содержимое диалога (спец 12.12).
|
// SetTranscript заменяет содержимое диалога (спец 12.12).
|
||||||
func (p *ChatPanel) SetTranscript(text string) {
|
func (p *ChatPanel) SetTranscript(text string) {
|
||||||
p.convLbl.SetText(text)
|
p.convLbl.SetText(text)
|
||||||
|
|||||||
@@ -116,21 +116,39 @@ func (p *LogPanel) flush() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// render перестраивает текст ленты из буфера с учётом выбранного набора
|
// render перерисовывает ленту: дописывает в label только новые видимые
|
||||||
// уровней. Пустой набор → пустая лента (без ошибок).
|
// записи (RenderNew). Если вытеснение задело уже выведенные записи (stale),
|
||||||
|
// RenderNew возвращает полный текст — пересобираем label целиком. Так при
|
||||||
|
// росте ленты мимо лимита флаш дешёвый (только дельта), а полная пересборка
|
||||||
|
// происходит лишь по вытеснению/смене фильтра.
|
||||||
func (p *LogPanel) render() {
|
func (p *LogPanel) render() {
|
||||||
p.label.SetText(p.buf.Render(p.selected))
|
delta, full := p.buf.RenderNew(p.selected)
|
||||||
|
if full {
|
||||||
|
p.label.SetText(delta)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if delta == "" {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
cur := p.label.Text
|
||||||
|
if cur == "" {
|
||||||
|
p.label.SetText(delta)
|
||||||
|
} else {
|
||||||
|
p.label.SetText(cur + delta)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// setLevel включает/выключает показ уровня и сразу обновляет ленту.
|
// setLevel включает/выключает показ уровня и сразу обновляет ленту.
|
||||||
// Вызывается только с потока Fyne (колбэки чекбоксов).
|
// Вызывается только с потока Fyne (колбэки чекбоксов). Смена набора уровней
|
||||||
|
// меняет видимость уже выведенных записей, поэтому — полный пересбор.
|
||||||
func (p *LogPanel) setLevel(level string, on bool) {
|
func (p *LogPanel) setLevel(level string, on bool) {
|
||||||
if p.selected[level] == on {
|
if p.selected[level] == on {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
p.selected[level] = on
|
p.selected[level] = on
|
||||||
p.btn.SetText(p.filterCaption())
|
p.btn.SetText(p.filterCaption())
|
||||||
p.render()
|
p.label.SetText(p.buf.Render(p.selected))
|
||||||
|
p.buf.ResetRendered()
|
||||||
}
|
}
|
||||||
|
|
||||||
// filterCaption — подпись кнопки-дропдауна: текущий выбранный набор уровней.
|
// filterCaption — подпись кнопки-дропдауна: текущий выбранный набор уровней.
|
||||||
|
|||||||
@@ -150,6 +150,8 @@ func (w *Window) addTab(sess *ui.Session) {
|
|||||||
panel.SetCommands(ui.NewCommands(func(text string) {
|
panel.SetCommands(ui.NewCommands(func(text string) {
|
||||||
w.submitText(sess.UserID, sess.Address, text)
|
w.submitText(sess.UserID, sess.Address, text)
|
||||||
}))
|
}))
|
||||||
|
// Кнопка «Перезапустить» активна только у вкладки с задачей (/retry N).
|
||||||
|
panel.SetBoundTask(sess.BoundTask)
|
||||||
|
|
||||||
content := container.NewBorder(nil, panel.Composer(), nil, nil, panel.Transcript())
|
content := container.NewBorder(nil, panel.Composer(), nil, nil, panel.Transcript())
|
||||||
tab := container.NewTabItem(sess.Title, content)
|
tab := container.NewTabItem(sess.Title, content)
|
||||||
@@ -309,6 +311,9 @@ func (w *Window) bindTaskSession(chatID string, taskID int64) {
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
s.BoundTask = taskID
|
s.BoundTask = taskID
|
||||||
|
if panel := w.panels[addr]; panel != nil {
|
||||||
|
panel.SetBoundTask(taskID)
|
||||||
|
}
|
||||||
if tab := w.tabForAddress(addr); tab != nil {
|
if tab := w.tabForAddress(addr); tab != nil {
|
||||||
tab.Text = "Задача #" + ui.Itoa(taskID)
|
tab.Text = "Задача #" + ui.Itoa(taskID)
|
||||||
w.tabs.Refresh()
|
w.tabs.Refresh()
|
||||||
|
|||||||
@@ -87,6 +87,13 @@ type LogBuffer struct {
|
|||||||
start int // индекс первого «живого» элемента (ленивое вытеснение)
|
start int // индекс первого «живого» элемента (ленивое вытеснение)
|
||||||
total int // суммарный объём текста живых записей [start..], байты
|
total int // суммарный объём текста живых записей [start..], байты
|
||||||
maxLen int // лимит суммарного объёма текста записей, байты
|
maxLen int // лимит суммарного объёма текста записей, байты
|
||||||
|
|
||||||
|
// Инкрементальная отрисовка (RenderNew): rendered — сколько живых записей
|
||||||
|
// от начала уже выведено в текст ленты; stale — вытеснение задело уже
|
||||||
|
// выведенные записи, так что текст ленты нельзя дополнить дельтом и нужен
|
||||||
|
// полный пересбор. Эти поля поддерживаются только рендером.
|
||||||
|
rendered int
|
||||||
|
stale bool
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewLogBuffer создаёт буфер ёмкостью maxLen байт текста записей.
|
// NewLogBuffer создаёт буфер ёмкостью maxLen байт текста записей.
|
||||||
@@ -113,18 +120,36 @@ func (b *LogBuffer) Append(level, text string) bool {
|
|||||||
func (b *LogBuffer) trim() bool {
|
func (b *LogBuffer) trim() bool {
|
||||||
dropped := false
|
dropped := false
|
||||||
live := len(b.entries) - b.start
|
live := len(b.entries) - b.start
|
||||||
|
evicted := 0
|
||||||
for b.total > b.maxLen && live > 1 {
|
for b.total > b.maxLen && live > 1 {
|
||||||
b.total -= len(b.entries[b.start].Text) + 1
|
b.total -= len(b.entries[b.start].Text) + 1
|
||||||
b.start++
|
b.start++
|
||||||
live--
|
live--
|
||||||
|
evicted++
|
||||||
dropped = true
|
dropped = true
|
||||||
}
|
}
|
||||||
|
// Если вытеснена хотя бы одна уже выведенная запись (b.rendered), текст
|
||||||
|
// ленты устарел: дельта не может убрать верхние строки, нужен полный
|
||||||
|
// пересбор. rendered при этом уменьшается на число вытесненных.
|
||||||
|
if evicted > 0 && b.rendered > 0 {
|
||||||
|
b.stale = true
|
||||||
|
}
|
||||||
|
if evicted >= b.rendered {
|
||||||
|
b.rendered = 0
|
||||||
|
} else {
|
||||||
|
b.rendered -= evicted
|
||||||
|
}
|
||||||
// Компакция «мёртвых» записей пачкой, а не на каждом append.
|
// Компакция «мёртвых» записей пачкой, а не на каждом append.
|
||||||
if b.start >= compactThreshold {
|
if b.start >= compactThreshold {
|
||||||
b.entries = b.entries[b.start:]
|
b.entries = b.entries[b.start:]
|
||||||
b.start = 0
|
b.start = 0
|
||||||
}
|
}
|
||||||
if n := len(b.entries); n > b.start && len(b.entries[n-1].Text) > b.maxLen {
|
if n := len(b.entries); n > b.start && len(b.entries[n-1].Text) > b.maxLen {
|
||||||
|
// Обрезка последней записи меняет уже выведенный текст, если она
|
||||||
|
// была отрисована.
|
||||||
|
if b.rendered == n-b.start {
|
||||||
|
b.stale = true
|
||||||
|
}
|
||||||
b.entries[n-1].Text = b.entries[n-1].Text[len(b.entries[n-1].Text)-b.maxLen:]
|
b.entries[n-1].Text = b.entries[n-1].Text[len(b.entries[n-1].Text)-b.maxLen:]
|
||||||
b.total = b.maxLen + 1
|
b.total = b.maxLen + 1
|
||||||
dropped = true
|
dropped = true
|
||||||
@@ -141,12 +166,55 @@ func (b *LogBuffer) Entries() []LogEntry {
|
|||||||
|
|
||||||
// Render строит текст ленты только из живых записей уровней, отмеченных в
|
// Render строит текст ленты только из живых записей уровней, отмеченных в
|
||||||
// selected (уровень → показывать). Пустой или nil набор → пустая лента.
|
// selected (уровень → показывать). Пустой или nil набор → пустая лента.
|
||||||
|
// Полный рендер: используется для первичной отрисовки и после смены фильтра.
|
||||||
func (b *LogBuffer) Render(selected map[string]bool) string {
|
func (b *LogBuffer) Render(selected map[string]bool) string {
|
||||||
if selected == nil {
|
if selected == nil {
|
||||||
return ""
|
return ""
|
||||||
}
|
}
|
||||||
|
return b.renderFrom(b.start, selected)
|
||||||
|
}
|
||||||
|
|
||||||
|
// ResetRendered помечает все живые записи как уже выведенные. Зовётся после
|
||||||
|
// полного рендера (SetText(Render)), чтобы следующий RenderNew вернул только
|
||||||
|
// новые записи.
|
||||||
|
func (b *LogBuffer) ResetRendered() {
|
||||||
|
b.rendered = len(b.entries) - b.start
|
||||||
|
b.stale = false
|
||||||
|
}
|
||||||
|
|
||||||
|
// RenderNew отдаёт текст ещё не выведенных видимых записей и флаг, требует
|
||||||
|
// ли слой полного пересбора ленты. Если вытеснение задело уже выведенные
|
||||||
|
// записи (stale), возвращаемый текст — полный рендер с нуля (и флаг true), а
|
||||||
|
// не дельта. Иначе — только текст новых записей, которым можно дополнить
|
||||||
|
// текущий текст ленты.
|
||||||
|
func (b *LogBuffer) RenderNew(selected map[string]bool) (text string, full bool) {
|
||||||
|
if selected == nil {
|
||||||
|
return "", false
|
||||||
|
}
|
||||||
|
live := len(b.entries) - b.start
|
||||||
|
if b.stale {
|
||||||
|
b.stale = false
|
||||||
|
b.rendered = live
|
||||||
|
return b.renderFrom(b.start, selected), true
|
||||||
|
}
|
||||||
|
from := b.start + b.rendered
|
||||||
var s strings.Builder
|
var s strings.Builder
|
||||||
for i := b.start; i < len(b.entries); i++ {
|
for i := from; i < len(b.entries); i++ {
|
||||||
|
e := b.entries[i]
|
||||||
|
if !selected[e.Level] {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
s.WriteString(e.Text)
|
||||||
|
s.WriteByte('\n')
|
||||||
|
}
|
||||||
|
b.rendered = live
|
||||||
|
return s.String(), false
|
||||||
|
}
|
||||||
|
|
||||||
|
// renderFrom строит текст записей, начиная с индекса from, с учётом selected.
|
||||||
|
func (b *LogBuffer) renderFrom(from int, selected map[string]bool) string {
|
||||||
|
var s strings.Builder
|
||||||
|
for i := from; i < len(b.entries); i++ {
|
||||||
e := b.entries[i]
|
e := b.entries[i]
|
||||||
if !selected[e.Level] {
|
if !selected[e.Level] {
|
||||||
continue
|
continue
|
||||||
|
|||||||
@@ -70,9 +70,9 @@ func allLevels() map[string]bool {
|
|||||||
func TestLogBufferRenderAllByDefault(t *testing.T) {
|
func TestLogBufferRenderAllByDefault(t *testing.T) {
|
||||||
b := NewLogBuffer(1000)
|
b := NewLogBuffer(1000)
|
||||||
b.Append("error", "e1")
|
b.Append("error", "e1")
|
||||||
b.Append("warn", "w1") // синоним → warning
|
b.Append("warn", "w1") // синоним → warning
|
||||||
b.Append("log", "i1") // неизвестный уровень → info
|
b.Append("log", "i1") // неизвестный уровень → info
|
||||||
b.Append("trace", "d1") // синоним → debug
|
b.Append("trace", "d1") // синоним → debug
|
||||||
|
|
||||||
got := b.Render(allLevels())
|
got := b.Render(allLevels())
|
||||||
want := "e1\nw1\ni1\nd1\n"
|
want := "e1\nw1\ni1\nd1\n"
|
||||||
@@ -229,3 +229,92 @@ func TestLogBufferBytesCap(t *testing.T) {
|
|||||||
t.Fatalf("oversized entry text = %q, want %q", entries[0].Text, "longtext")
|
t.Fatalf("oversized entry text = %q, want %q", entries[0].Text, "longtext")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Первый RenderNew отдаёт весь видимый текст, последующие — только новые
|
||||||
|
// записи (инкрементальная отрисовка: повторный вызов без Append пуст).
|
||||||
|
func TestLogBufferRenderNewDelta(t *testing.T) {
|
||||||
|
b := NewLogBuffer(1000)
|
||||||
|
sel := allLevels()
|
||||||
|
|
||||||
|
b.Append(LogLevelError, "e1")
|
||||||
|
b.Append(LogLevelInfo, "i1")
|
||||||
|
text, full := b.RenderNew(sel)
|
||||||
|
if full {
|
||||||
|
t.Fatalf("первый RenderNew не должен требовать полного пересбора, got full")
|
||||||
|
}
|
||||||
|
if text != "e1\ni1\n" {
|
||||||
|
t.Fatalf("первый RenderNew = %q, want %q", text, "e1\ni1\n")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Без новых записей — пустой дельта.
|
||||||
|
text, full = b.RenderNew(sel)
|
||||||
|
if full || text != "" {
|
||||||
|
t.Fatalf("повторный RenderNew без новых записей = (%q, %v), want (\"\", false)", text, full)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Новые записи — только их текст, без повторной отдачи старых.
|
||||||
|
b.Append(LogLevelInfo, "i2")
|
||||||
|
text, _ = b.RenderNew(sel)
|
||||||
|
if text != "i2\n" {
|
||||||
|
t.Fatalf("дельта RenderNew = %q, want %q", text, "i2\n")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// RenderNew применяет фильтр: невидимые уровни не попадают в дельту.
|
||||||
|
func TestLogBufferRenderNewFiltered(t *testing.T) {
|
||||||
|
b := NewLogBuffer(1000)
|
||||||
|
sel := map[string]bool{LogLevelError: true}
|
||||||
|
b.Append(LogLevelError, "e1")
|
||||||
|
b.Append(LogLevelInfo, "i1")
|
||||||
|
b.Append(LogLevelError, "e2")
|
||||||
|
|
||||||
|
text, _ := b.RenderNew(sel)
|
||||||
|
if text != "e1\ne2\n" {
|
||||||
|
t.Fatalf("RenderNew(error only) = %q, want %q", text, "e1\ne2\n")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// RenderNew уважает ResetRendered: после полного рендера новые записи
|
||||||
|
// добавляются дельтом, а не вытесняют уже выведенный текст.
|
||||||
|
func TestLogBufferRenderNewAfterReset(t *testing.T) {
|
||||||
|
b := NewLogBuffer(1000)
|
||||||
|
sel := allLevels()
|
||||||
|
b.Append(LogLevelInfo, "old")
|
||||||
|
b.Append(LogLevelInfo, "base")
|
||||||
|
|
||||||
|
text, full := b.RenderNew(sel)
|
||||||
|
if full || text != "old\nbase\n" {
|
||||||
|
t.Fatalf("RenderNew = (%q, %v), want (\"old\\nbase\\n\", false)", text, full)
|
||||||
|
}
|
||||||
|
b.ResetRendered() // имитация полной пересборки ленты
|
||||||
|
|
||||||
|
b.Append(LogLevelInfo, "new")
|
||||||
|
text, full = b.RenderNew(sel)
|
||||||
|
if full || text != "new\n" {
|
||||||
|
t.Fatalf("RenderNew после Reset = (%q, %v), want (\"new\\n\", false)", text, full)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Вытеснение уже выведенных записей помечает рендер как требующий полного
|
||||||
|
// пересбора (stale): дельта не может убрать верхние строки.
|
||||||
|
func TestLogBufferRenderNewEvictionStale(t *testing.T) {
|
||||||
|
b := NewLogBuffer(25) // записи по 10 байт + \n: помещаются 2, 3-я вытесняет 1-ю
|
||||||
|
sel := allLevels()
|
||||||
|
|
||||||
|
b.Append(LogLevelInfo, "0123456789")
|
||||||
|
b.Append(LogLevelInfo, "0123456789")
|
||||||
|
text, full := b.RenderNew(sel)
|
||||||
|
if full || text != "0123456789\n0123456789\n" {
|
||||||
|
t.Fatalf("RenderNew = (%q, %v), want two lines", text, full)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Третья запись вытесняет первую (уже отрисованную) — stale.
|
||||||
|
b.Append(LogLevelInfo, "0123456789")
|
||||||
|
text, full = b.RenderNew(sel)
|
||||||
|
if !full {
|
||||||
|
t.Fatalf("RenderNew после вытеснения = (full=%v), want full=true", full)
|
||||||
|
}
|
||||||
|
if text != "0123456789\n0123456789\n" {
|
||||||
|
t.Fatalf("полный текст после вытеснения = %q, want последние две записи", text)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -31,7 +31,9 @@ import (
|
|||||||
"os/exec"
|
"os/exec"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"runtime"
|
"runtime"
|
||||||
|
"sort"
|
||||||
"strings"
|
"strings"
|
||||||
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -50,6 +52,11 @@ type Updater struct {
|
|||||||
// Dir — каталог рядом с бинарём (для .new/.old и companion-метаданных).
|
// Dir — каталог рядом с бинарём (для .new/.old и companion-метаданных).
|
||||||
// Ставится app из os.Executable(); если пуст — используется каталог Workdir.
|
// Ставится app из os.Executable(); если пуст — используется каталог Workdir.
|
||||||
Dir string
|
Dir string
|
||||||
|
|
||||||
|
// client — общий HTTP-клиент (keep-alive), чтобы проверки/скачивания
|
||||||
|
// переиспользовали соединения, а не создавали новое на каждый запрос.
|
||||||
|
client *http.Client
|
||||||
|
clientMu sync.Mutex
|
||||||
}
|
}
|
||||||
|
|
||||||
// Result — результат Check.
|
// Result — результат Check.
|
||||||
@@ -118,16 +125,26 @@ func (u *Updater) versionsURL() string {
|
|||||||
return base + "/api/v1/packages/" + url.PathEscape(u.Owner) + "/generic/" + url.PathEscape(u.Package)
|
return base + "/api/v1/packages/" + url.PathEscape(u.Owner) + "/generic/" + url.PathEscape(u.Package)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// httpClient возвращает общий клиент (keep-alive), инициализируя лениво.
|
||||||
|
func (u *Updater) httpClient() *http.Client {
|
||||||
|
u.clientMu.Lock()
|
||||||
|
defer u.clientMu.Unlock()
|
||||||
|
if u.client == nil {
|
||||||
|
u.client = &http.Client{Timeout: 30 * time.Second}
|
||||||
|
}
|
||||||
|
return u.client
|
||||||
|
}
|
||||||
|
|
||||||
// httpGet скачивает файл по URL бэкенда. При Token непустом — Basic/токен-заголовок.
|
// httpGet скачивает файл по URL бэкенда. При Token непустом — Basic/токен-заголовок.
|
||||||
func (u *Updater) httpGet(url string) ([]byte, error) {
|
func (u *Updater) httpGet(ctx context.Context, url string) ([]byte, error) {
|
||||||
req, err := http.NewRequest(http.MethodGet, url, nil)
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
if u.Token != "" {
|
if u.Token != "" {
|
||||||
req.Header.Set("Authorization", "token "+u.Token)
|
req.Header.Set("Authorization", "token "+u.Token)
|
||||||
}
|
}
|
||||||
resp, err := (&http.Client{Timeout: 30 * time.Second}).Do(req)
|
resp, err := u.httpClient().Do(req)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
@@ -138,6 +155,61 @@ func (u *Updater) httpGet(url string) ([]byte, error) {
|
|||||||
return io.ReadAll(resp.Body)
|
return io.ReadAll(resp.Body)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// fileExists проверяет наличие файла по URL без скачивания тела: HEAD,
|
||||||
|
// а при 405/501 (сервер не поддерживает HEAD) — fallback на GET с Range байт 0-0.
|
||||||
|
// Возвращает (false, nil) при 404/410 — файла нет.
|
||||||
|
func (u *Updater) fileExists(ctx context.Context, url string) (bool, error) {
|
||||||
|
req, err := http.NewRequestWithContext(ctx, http.MethodHead, url, nil)
|
||||||
|
if err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
|
if u.Token != "" {
|
||||||
|
req.Header.Set("Authorization", "token "+u.Token)
|
||||||
|
}
|
||||||
|
resp, err := u.httpClient().Do(req)
|
||||||
|
if err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
|
resp.Body.Close()
|
||||||
|
switch {
|
||||||
|
case resp.StatusCode == http.StatusNotFound || resp.StatusCode == http.StatusGone:
|
||||||
|
return false, nil
|
||||||
|
case resp.StatusCode == http.StatusOK || resp.StatusCode == http.StatusPartialContent:
|
||||||
|
return true, nil
|
||||||
|
case resp.StatusCode == http.StatusMethodNotAllowed || resp.StatusCode == http.StatusNotImplemented:
|
||||||
|
// Gitea может не отвечать на HEAD — проверяем GET с Range 0-0 без чтения тела.
|
||||||
|
return u.fileExistsByRange(ctx, url)
|
||||||
|
default:
|
||||||
|
return false, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// fileExistsByRange проверяет наличие файла GET'ом с Range: bytes=0-0.
|
||||||
|
// Тело не читается: достаточно лишь первых байт заголовков ответа.
|
||||||
|
func (u *Updater) fileExistsByRange(ctx context.Context, url string) (bool, error) {
|
||||||
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
||||||
|
if err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
|
if u.Token != "" {
|
||||||
|
req.Header.Set("Authorization", "token "+u.Token)
|
||||||
|
}
|
||||||
|
req.Header.Set("Range", "bytes=0-0")
|
||||||
|
resp, err := u.httpClient().Do(req)
|
||||||
|
if err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
|
resp.Body.Close()
|
||||||
|
switch {
|
||||||
|
case resp.StatusCode == http.StatusNotFound || resp.StatusCode == http.StatusGone:
|
||||||
|
return false, nil
|
||||||
|
case resp.StatusCode == http.StatusOK || resp.StatusCode == http.StatusPartialContent:
|
||||||
|
return true, nil
|
||||||
|
default:
|
||||||
|
return false, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Check определяет, есть ли более свежая версия в Gitea Packages.
|
// Check определяет, есть ли более свежая версия в Gitea Packages.
|
||||||
// Не скачивает бинарь. Ошибка (U1) возвращается в Result.Err — вызывающий
|
// Не скачивает бинарь. Ошибка (U1) возвращается в Result.Err — вызывающий
|
||||||
// решает, логировать и пропустить.
|
// решает, логировать и пропустить.
|
||||||
@@ -167,8 +239,12 @@ type pkgVersion struct {
|
|||||||
// ResolveLatest определяет идентификатор новейшей применимой версии пакета.
|
// ResolveLatest определяет идентификатор новейшей применимой версии пакета.
|
||||||
// Бинарь/метаданные читаем из КОНКРЕТНОЙ версии, а не из pseudo-`latest`,
|
// Бинарь/метаданные читаем из КОНКРЕТНОЙ версии, а не из pseudo-`latest`,
|
||||||
// чтобы companion-файлы и бинарь всегда брались из одного снимка.
|
// чтобы companion-файлы и бинарь всегда брались из одного снимка.
|
||||||
|
//
|
||||||
|
// Применимость версии проверяем НАЛИЧИЕМ бинаря платформы (HEAD без тела),
|
||||||
|
// а не скачиванием полного файла: при N версиях это N запросов заголовков
|
||||||
|
// вместо N×(размер бинаря) байт.
|
||||||
func (u *Updater) ResolveLatest(ctx context.Context) (string, error) {
|
func (u *Updater) ResolveLatest(ctx context.Context) (string, error) {
|
||||||
b, err := u.httpGet(u.versionsURL())
|
b, err := u.httpGet(ctx, u.versionsURL())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
@@ -176,30 +252,28 @@ func (u *Updater) ResolveLatest(ctx context.Context) (string, error) {
|
|||||||
if err := json.Unmarshal(b, &vers); err != nil {
|
if err := json.Unmarshal(b, &vers); err != nil {
|
||||||
return "", ue(U1, "list "+u.versionsURL(), err)
|
return "", ue(U1, "list "+u.versionsURL(), err)
|
||||||
}
|
}
|
||||||
|
// новые версии — с большим ID; идём с новейшей и берём первую с бинарём.
|
||||||
|
sort.SliceStable(vers, func(i, j int) bool { return vers[i].ID > vers[j].ID })
|
||||||
want := PlatformFilename()
|
want := PlatformFilename()
|
||||||
// выбираем самую свежую версию (макс. id) класса commit-*, в которой есть бинарь.
|
|
||||||
best := ""
|
|
||||||
var bestID int64
|
|
||||||
for _, v := range vers {
|
for _, v := range vers {
|
||||||
if !strings.HasPrefix(v.Ver, "commit-") {
|
if !strings.HasPrefix(v.Ver, "commit-") {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
// версия применима, только если в ней опубликован бинарь нашей платформы
|
exists, ferr := u.fileExists(ctx, u.fileURL(v.Ver, want))
|
||||||
if _, err := u.httpGet(u.fileURL(v.Ver, want)); err != nil {
|
if ferr != nil {
|
||||||
continue
|
continue // сетевые ошибки пробы не роняют проверку
|
||||||
}
|
}
|
||||||
if v.ID > bestID {
|
if exists {
|
||||||
bestID = v.ID
|
return v.Ver, nil
|
||||||
best = v.Ver
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return best, nil
|
return "", nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Download скачивает бинарь конкретной версии во временный файл и возвращает
|
// Download скачивает бинарь конкретной версии во временный файл и возвращает
|
||||||
// путь к нему. Файл: <Dir>/.ratatoskr.<ver>.new.
|
// путь к нему. Файл: <Dir>/.ratatoskr.<ver>.new.
|
||||||
func (u *Updater) Download(ctx context.Context, version string) (string, error) {
|
func (u *Updater) Download(ctx context.Context, version string) (string, error) {
|
||||||
b, err := u.httpGet(u.fileURL(version, PlatformFilename()))
|
b, err := u.httpGet(ctx, u.fileURL(version, PlatformFilename()))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err // уже U3
|
return "", err // уже U3
|
||||||
}
|
}
|
||||||
@@ -216,7 +290,7 @@ func (u *Updater) Download(ctx context.Context, version string) (string, error)
|
|||||||
|
|
||||||
// versionSum256 читает companion-файл контрольной суммы конкретной версии.
|
// versionSum256 читает companion-файл контрольной суммы конкретной версии.
|
||||||
func (u *Updater) versionSum256(ctx context.Context, version string) (string, error) {
|
func (u *Updater) versionSum256(ctx context.Context, version string) (string, error) {
|
||||||
b, err := u.httpGet(u.fileURL(version, PlatformFilename()+".sha256"))
|
b, err := u.httpGet(ctx, u.fileURL(version, PlatformFilename()+".sha256"))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
@@ -352,4 +426,4 @@ func (u *Updater) SwapAndRestart(file string) error {
|
|||||||
// успешный старт нового процесса — текущий завершаем
|
// успешный старт нового процесса — текущий завершаем
|
||||||
os.Exit(0)
|
os.Exit(0)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -9,6 +9,7 @@ import (
|
|||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
|
"sync/atomic"
|
||||||
"testing"
|
"testing"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -41,7 +42,7 @@ func mockGitea(t *testing.T, bin []byte, version, checksum string) *httptest.Ser
|
|||||||
_ = json.NewEncoder(w).Encode(versions)
|
_ = json.NewEncoder(w).Encode(versions)
|
||||||
})
|
})
|
||||||
mux.HandleFunc("/api/packages/", func(w http.ResponseWriter, r *http.Request) {
|
mux.HandleFunc("/api/packages/", func(w http.ResponseWriter, r *http.Request) {
|
||||||
if r.Method != http.MethodGet {
|
if r.Method != http.MethodGet && r.Method != http.MethodHead {
|
||||||
http.Error(w, "method", http.StatusMethodNotAllowed)
|
http.Error(w, "method", http.StatusMethodNotAllowed)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -58,7 +59,10 @@ func mockGitea(t *testing.T, bin []byte, version, checksum string) *httptest.Ser
|
|||||||
http.NotFound(w, r)
|
http.NotFound(w, r)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
_, _ = w.Write(body)
|
if r.Method == http.MethodGet {
|
||||||
|
_, _ = w.Write(body)
|
||||||
|
}
|
||||||
|
// HEAD — просто 200, тело не пишем
|
||||||
})
|
})
|
||||||
return httptest.NewServer(mux)
|
return httptest.NewServer(mux)
|
||||||
}
|
}
|
||||||
@@ -133,6 +137,104 @@ func TestCheck_ServerDown(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestResolveLatest_HeadNotFullDownload проверяет, что ResolveLatest проверяет
|
||||||
|
// наличие бинаря HEAD'ом и НЕ скачивает полное тело бинаря (прошлая версия
|
||||||
|
// читала каждый файл целиком — O(N)×размер бинаря).
|
||||||
|
func TestResolveLatest_HeadNotFullDownload(t *testing.T) {
|
||||||
|
name := PlatformFilename()
|
||||||
|
_ = name
|
||||||
|
var headReqs, bodyReqs int64
|
||||||
|
mux := http.NewServeMux()
|
||||||
|
mux.HandleFunc("/api/v1/packages/", func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
_ = json.NewEncoder(w).Encode([]map[string]any{
|
||||||
|
{"id": 1, "version": "commit-aaa1111"},
|
||||||
|
{"id": 2, "version": "commit-abc1234"},
|
||||||
|
})
|
||||||
|
})
|
||||||
|
mux.HandleFunc("/api/packages/", func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
parts := strings.Split(strings.Trim(r.URL.Path, "/"), "/")
|
||||||
|
if len(parts) < 6 {
|
||||||
|
http.NotFound(w, r)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
ver := parts[len(parts)-2]
|
||||||
|
fileName := parts[len(parts)-1]
|
||||||
|
if ver != "commit-abc1234" || fileName != name {
|
||||||
|
http.NotFound(w, r)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
switch r.Method {
|
||||||
|
case http.MethodHead:
|
||||||
|
atomic.AddInt64(&headReqs, 1)
|
||||||
|
case http.MethodGet:
|
||||||
|
atomic.AddInt64(&bodyReqs, 1)
|
||||||
|
default:
|
||||||
|
http.Error(w, "method", http.StatusMethodNotAllowed)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
srv := httptest.NewServer(mux)
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
u := &Updater{BaseURL: srv.URL, Owner: "k", Package: "p", CurrentVersion: "v1", Dir: t.TempDir()}
|
||||||
|
ver, err := u.ResolveLatest(context.Background())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("ResolveLatest err = %v", err)
|
||||||
|
}
|
||||||
|
if ver != "commit-abc1234" {
|
||||||
|
t.Errorf("ResolveLatest = %q, want commit-abc1234", ver)
|
||||||
|
}
|
||||||
|
if atomic.LoadInt64(&headReqs) == 0 {
|
||||||
|
t.Error("ResolveLatest не делал HEAD-проб на файлы")
|
||||||
|
}
|
||||||
|
if atomic.LoadInt64(&bodyReqs) != 0 {
|
||||||
|
t.Errorf("ResolveLatest скачал тело бинарника: %d полных GET", atomic.LoadInt64(&bodyReqs))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestResolveLatest_SkipsBinarylessVersion проверяет фоллбэк: новейшая версия
|
||||||
|
// без бинаря (разные матрицы платформ публикуются не все сразу) пропускается,
|
||||||
|
// берётся следующая, где файл есть.
|
||||||
|
func TestResolveLatest_SkipsBinarylessVersion(t *testing.T) {
|
||||||
|
name := PlatformFilename()
|
||||||
|
bin := []byte("binary")
|
||||||
|
mux := http.NewServeMux()
|
||||||
|
mux.HandleFunc("/api/v1/packages/", func(w http.ResponseWriter, _ *http.Request) {
|
||||||
|
_ = json.NewEncoder(w).Encode([]map[string]any{
|
||||||
|
{"id": 1, "version": "commit-aaa1111"},
|
||||||
|
{"id": 2, "version": "commit-mid2222"},
|
||||||
|
{"id": 3, "version": "commit-new3333"},
|
||||||
|
})
|
||||||
|
})
|
||||||
|
mux.HandleFunc("/api/packages/", func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
parts := strings.Split(strings.Trim(r.URL.Path, "/"), "/")
|
||||||
|
if len(parts) < 6 {
|
||||||
|
http.NotFound(w, r)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
ver := parts[len(parts)-2]
|
||||||
|
fileName := parts[len(parts)-1]
|
||||||
|
// бинарь есть только у commit-mid2222 — новейшие 1 и 3 пропускаются
|
||||||
|
if ver != "commit-mid2222" || fileName != name {
|
||||||
|
http.NotFound(w, r)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if r.Method != http.MethodHead {
|
||||||
|
_, _ = w.Write(bin)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
srv := httptest.NewServer(mux)
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
u := &Updater{BaseURL: srv.URL, Owner: "k", Package: "p", CurrentVersion: "v0", Dir: t.TempDir()}
|
||||||
|
ver, err := u.ResolveLatest(context.Background())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("ResolveLatest err = %v", err)
|
||||||
|
}
|
||||||
|
if ver != "commit-mid2222" {
|
||||||
|
t.Errorf("ResolveLatest = %q, want commit-mid2222 (фоллбэк от версии без бинаря)", ver)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestDownload_And_Verify_Good(t *testing.T) {
|
func TestDownload_And_Verify_Good(t *testing.T) {
|
||||||
bin := []byte("ratatoskr-binary-content-v2")
|
bin := []byte("ratatoskr-binary-content-v2")
|
||||||
u, dir := testUpdater(t, bin, "commit-new12345", "")
|
u, dir := testUpdater(t, bin, "commit-new12345", "")
|
||||||
@@ -249,4 +351,4 @@ func errorsAs(err error, target **Error) bool {
|
|||||||
err = c.Unwrap()
|
err = c.Unwrap()
|
||||||
}
|
}
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|||||||
205
internal/worker/postmortem.go
Normal file
205
internal/worker/postmortem.go
Normal file
@@ -0,0 +1,205 @@
|
|||||||
|
package worker
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"log"
|
||||||
|
"strings"
|
||||||
|
"text/template"
|
||||||
|
|
||||||
|
"github.com/kamelion/ratatoskr-go/internal/events"
|
||||||
|
"github.com/kamelion/ratatoskr-go/internal/storage"
|
||||||
|
)
|
||||||
|
|
||||||
|
// postMortemAgent — имя постмортем-агента (файл agents/postmortem.md).
|
||||||
|
const postMortemAgent = "postmortem"
|
||||||
|
|
||||||
|
// traceOutputMax — обрезка вывода сессии в постмортем-промпте (чтобы промпт
|
||||||
|
// не превращался в полные транскрипты и не переполнял контекст модели).
|
||||||
|
const traceOutputMax = 6000
|
||||||
|
|
||||||
|
// postMortemPromptTemplate — промпт для постмортем-агента после failed/timeout:
|
||||||
|
// задача + сессии dev/reviewer. Ожидается резюме простым текстом на русском.
|
||||||
|
var postMortemPromptTemplate = template.Must(template.New("postmortem").Parse(`Ты — постмортем-аналитик в конвейере Ratatoskr. Задача завершилась неудачей ({{.Status}}). Проанализируй сессии агентов dev/reviewer и дай резюме: почему так случилось и что сделать, чтобы не повторялось.
|
||||||
|
|
||||||
|
**Задача:**
|
||||||
|
{{if .Title}}Название: {{.Title}}{{end}}
|
||||||
|
{{if .Goal}}Цель: {{.Goal}}{{end}}
|
||||||
|
{{if .Repos}}
|
||||||
|
Репозитории:
|
||||||
|
{{- range .Repos}}
|
||||||
|
- {{.}}
|
||||||
|
{{- end}}
|
||||||
|
{{end}}
|
||||||
|
{{if .Why}}Зачем: {{.Why}}{{end}}
|
||||||
|
{{if .AC}}Критерии готовности (AC):
|
||||||
|
{{.AC}}{{end}}
|
||||||
|
|
||||||
|
**Итоговый статус задачи:** {{.Status}}
|
||||||
|
|
||||||
|
**Сессии субагентов:**
|
||||||
|
{{.Sessions}}
|
||||||
|
|
||||||
|
Ответь ПРОСТЫМ ТЕКСТОМ на русском, без JSON и разметки. Формат:
|
||||||
|
|
||||||
|
Почему так случилось:
|
||||||
|
- <причина 1>
|
||||||
|
- <причина 2>
|
||||||
|
|
||||||
|
Что сделать, чтобы это не повторялось:
|
||||||
|
- <рекомендация 1>
|
||||||
|
- <рекомендация 2>`))
|
||||||
|
|
||||||
|
// PostMortemPromptData — данные для рендера постмортем-промпта.
|
||||||
|
type PostMortemPromptData struct {
|
||||||
|
Title string
|
||||||
|
Goal string
|
||||||
|
Repos []string
|
||||||
|
Why string
|
||||||
|
AC string
|
||||||
|
Status storage.Status
|
||||||
|
Sessions string
|
||||||
|
}
|
||||||
|
|
||||||
|
// RenderPostMortemPrompt собирает промпт для постмортем-агента.
|
||||||
|
func RenderPostMortemPrompt(data PostMortemPromptData) (string, error) {
|
||||||
|
var buf strings.Builder
|
||||||
|
if err := postMortemPromptTemplate.Execute(&buf, data); err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
return buf.String(), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// postMortemsText форматирует сессии dev/reviewer в секцию промпта.
|
||||||
|
func postMortemsText(traces []storage.Trace) string {
|
||||||
|
var b strings.Builder
|
||||||
|
for _, tr := range traces {
|
||||||
|
b.WriteString("\n=== Агент: " + tr.Agent + " (статус " + string(tr.Status) + ") ===\n")
|
||||||
|
if tr.SessionID != "" {
|
||||||
|
b.WriteString("session_id: " + tr.SessionID + "\n")
|
||||||
|
}
|
||||||
|
if strings.TrimSpace(tr.Prompt) != "" {
|
||||||
|
b.WriteString("-- Промпт агента --\n")
|
||||||
|
b.WriteString(tr.Prompt)
|
||||||
|
b.WriteString("\n")
|
||||||
|
}
|
||||||
|
if strings.TrimSpace(tr.Output) != "" {
|
||||||
|
b.WriteString("-- Вывод агента --\n")
|
||||||
|
b.WriteString(truncateTrace(tr.Output, traceOutputMax))
|
||||||
|
b.WriteString("\n")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if b.Len() == 0 {
|
||||||
|
return "(сессии dev/reviewer не найдены — вероятна инфраструктурная ошибка до запуска агентов)"
|
||||||
|
}
|
||||||
|
return b.String()
|
||||||
|
}
|
||||||
|
|
||||||
|
// truncateTrace обрезает длинный текст до последних n символов (релевантен
|
||||||
|
// хвост: вердикт/ошибка агента в конце вывода).
|
||||||
|
func truncateTrace(s string, n int) string {
|
||||||
|
if len(s) <= n {
|
||||||
|
return s
|
||||||
|
}
|
||||||
|
if n <= 0 {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
return "(вывод обрезан)\n" + s[len(s)-n:]
|
||||||
|
}
|
||||||
|
|
||||||
|
// hasPostMortemTrace возвращает true, если у задачи уже есть постмортем-trace
|
||||||
|
// (защита от повторного запуска при повторных прогонах/retry).
|
||||||
|
func (w *Worker) hasPostMortemTrace(ctx context.Context, taskID int64) bool {
|
||||||
|
if w.Store == nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
_, err := w.Store.GetLatestTrace(ctx, taskID, postMortemAgent)
|
||||||
|
return err == nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// postMortem запускает постмортем-агент для задачи, завершившейся failed/timeout:
|
||||||
|
// собирает сессии dev/reviewer из трасс, даёт агенту анализ, сохраняет результат
|
||||||
|
// как trace agent=postmortem и уведомляет владельца задачи резюме.
|
||||||
|
//
|
||||||
|
// Статус задачи НЕ меняется (failed/timeout остаётся достигнутым); собственные
|
||||||
|
// сбои постмортема не влияют на исход задачи — только логируются.
|
||||||
|
func (w *Worker) postMortem(ctx context.Context, task *storage.Task) {
|
||||||
|
if w.Store == nil || w.Runner == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if w.hasPostMortemTrace(ctx, task.ID) {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
traces, err := w.Store.GetTraces(ctx, task.ID)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("worker: task %d: постмортем: трассы: %v", task.ID, err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
// Анализируем только сессии агентов конвейера (dev/reviewer).
|
||||||
|
var sessions []storage.Trace
|
||||||
|
for _, tr := range traces {
|
||||||
|
if tr.Agent == "dev" || tr.Agent == "reviewer" {
|
||||||
|
sessions = append(sessions, *tr)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
prompt, pErr := RenderPostMortemPrompt(PostMortemPromptData{
|
||||||
|
Title: task.Title,
|
||||||
|
Goal: task.Goal,
|
||||||
|
Repos: task.EffectiveRepos(),
|
||||||
|
Why: task.Why,
|
||||||
|
AC: task.AC,
|
||||||
|
Status: task.Status,
|
||||||
|
Sessions: postMortemsText(sessions),
|
||||||
|
})
|
||||||
|
if pErr != nil {
|
||||||
|
log.Printf("worker: task %d: постмортем: рендер промпта: %v", task.ID, pErr)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
// лог-событие для UI-панели «Состояние».
|
||||||
|
w.publish(events.AgentActivity{TaskID: task.ID, Agent: postMortemAgent, Stage: "postmortem"})
|
||||||
|
|
||||||
|
tr := &storage.Trace{TaskID: task.ID, Agent: postMortemAgent, Prompt: prompt}
|
||||||
|
traceID, aErr := w.Store.AppendTrace(ctx, tr)
|
||||||
|
if aErr != nil {
|
||||||
|
log.Printf("worker: task %d: постмортем: create trace: %v", task.ID, aErr)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
res, rErr := w.Runner.Run(w.runCtx(ctx, task.ID), prompt, w.Worktree, postMortemAgent, "")
|
||||||
|
if rErr != nil {
|
||||||
|
log.Printf("worker: task %d: постмортем: запуск: %v", task.ID, rErr)
|
||||||
|
w.finalizeTrace(ctx, traceID, storage.TraceFailed, rErr.Error())
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if res.SessionID != "" {
|
||||||
|
_ = w.Store.UpdateTraceSessionID(ctx, traceID, res.SessionID)
|
||||||
|
}
|
||||||
|
|
||||||
|
output := strings.TrimSpace(res.Stdout)
|
||||||
|
status := storage.TraceSuccess
|
||||||
|
if res.RC != 0 || output == "" {
|
||||||
|
status = storage.TraceFailed
|
||||||
|
if output == "" {
|
||||||
|
output = "(постмортем-агент не вернул текст)"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
w.finalizeTrace(ctx, traceID, status, output)
|
||||||
|
|
||||||
|
if res.RC == 0 && output != "" {
|
||||||
|
text := fmt.Sprintf("Задача #%d: 🔍 анализ (после %s)\n%s", task.ID, task.Status, output)
|
||||||
|
w.notify(ctx, task, text)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// postMortemAfter — defer-хук из runTask: запускает постмортем, если задача
|
||||||
|
// завершилась failed/timeout. Собственные ошибки постмортема не мешают
|
||||||
|
// исходному результату задачи (возвращаемый *error только читается).
|
||||||
|
func (w *Worker) postMortemAfter(ctx context.Context, task *storage.Task, _ *error) {
|
||||||
|
if task.Status != storage.StatusFailed && task.Status != storage.StatusTimeout {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
w.postMortem(ctx, task)
|
||||||
|
}
|
||||||
@@ -196,6 +196,10 @@ func (w *Worker) pollAndDispatch(ctx context.Context) error {
|
|||||||
// dev-агент реализует, reviewer строго проверяет весь дифф ветки; при не-проходе
|
// dev-агент реализует, reviewer строго проверяет весь дифф ветки; при не-проходе
|
||||||
// dev дорабатывает по комментариям; прошло → push ветки + success.
|
// dev дорабатывает по комментариям; прошло → push ветки + success.
|
||||||
func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
|
||||||
|
// 0. постмортем-анализ: если задача завершилась failed/timeout — после
|
||||||
|
// выхода из runTask (всех путей) запускаем анализ сессий dev/reviewer.
|
||||||
|
defer w.postMortemAfter(ctx, task, &err)
|
||||||
|
|
||||||
// 1. проверяем статус
|
// 1. проверяем статус
|
||||||
if task.Status != storage.StatusApproved {
|
if task.Status != storage.StatusApproved {
|
||||||
return fmt.Errorf("%w: task %d status=%q", ErrLaunch, task.ID, task.Status)
|
return fmt.Errorf("%w: task %d status=%q", ErrLaunch, task.ID, task.Status)
|
||||||
|
|||||||
@@ -54,6 +54,11 @@ type mockRunnerWorker struct {
|
|||||||
// (последний вечный). Позволяет смоделировать fail-then-pass.
|
// (последний вечный). Позволяет смоделировать fail-then-pass.
|
||||||
reviewSequence []*opencode.Result
|
reviewSequence []*opencode.Result
|
||||||
reviewCount int
|
reviewCount int
|
||||||
|
|
||||||
|
// postMortemResult — результат постмортем-агента; nil → RC:0 с тестовым
|
||||||
|
// резюме. postMortemCount — сколько раз постмортем вызывался.
|
||||||
|
postMortemResult *opencode.Result
|
||||||
|
postMortemCount int
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *mockRunnerWorker) Run(_ context.Context, _, _, agent, _ string) (*opencode.Result, error) {
|
func (m *mockRunnerWorker) Run(_ context.Context, _, _, agent, _ string) (*opencode.Result, error) {
|
||||||
@@ -71,6 +76,13 @@ func (m *mockRunnerWorker) Run(_ context.Context, _, _, agent, _ string) (*openc
|
|||||||
}
|
}
|
||||||
return &opencode.Result{RC: 0, Stdout: `{"passed":true,"comments":[]}`}, nil
|
return &opencode.Result{RC: 0, Stdout: `{"passed":true,"comments":[]}`}, nil
|
||||||
}
|
}
|
||||||
|
if agent == postMortemAgent {
|
||||||
|
m.postMortemCount++
|
||||||
|
if m.postMortemResult != nil {
|
||||||
|
return m.postMortemResult, m.err
|
||||||
|
}
|
||||||
|
return &opencode.Result{RC: 0, Stdout: "Почему так случилось:\n- тестовое резюме\nЧто сделать:\n- исправить"}, nil
|
||||||
|
}
|
||||||
return m.result, m.err
|
return m.result, m.err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -402,12 +414,17 @@ func TestWorkerIterationsLimitNotification(t *testing.T) {
|
|||||||
if len(texts) == 0 {
|
if len(texts) == 0 {
|
||||||
t.Fatal("нет уведомлений")
|
t.Fatal("нет уведомлений")
|
||||||
}
|
}
|
||||||
last := texts[len(texts)-1]
|
// где-то в списке есть уведомление о failed с упоминанием лимита итераций
|
||||||
if !strings.Contains(last, "failed") {
|
// (над ним — уведомление postmortem, поэтому последним оно не обязано быть).
|
||||||
t.Errorf("последнее уведомление = %q, want упоминание failed", last)
|
var limitNotif string
|
||||||
|
for _, txt := range texts {
|
||||||
|
if strings.Contains(txt, "failed") && strings.Contains(txt, "итераци") {
|
||||||
|
limitNotif = txt
|
||||||
|
break
|
||||||
|
}
|
||||||
}
|
}
|
||||||
if !strings.Contains(last, "итераци") {
|
if limitNotif == "" {
|
||||||
t.Errorf("последнее уведомление = %q, want упоминание лимита итераций", last)
|
t.Errorf("нет уведомления о лимите итераций (failed+итераци): %#v", texts)
|
||||||
}
|
}
|
||||||
// ровно одно уведомление о failed (running→failed не задваивается)
|
// ровно одно уведомление о failed (running→failed не задваивается)
|
||||||
var failedCount int
|
var failedCount int
|
||||||
@@ -419,9 +436,20 @@ func TestWorkerIterationsLimitNotification(t *testing.T) {
|
|||||||
if failedCount != 1 {
|
if failedCount != 1 {
|
||||||
t.Errorf("уведомлений о failed = %d, want ровно 1: %#v", failedCount, texts)
|
t.Errorf("уведомлений о failed = %d, want ровно 1: %#v", failedCount, texts)
|
||||||
}
|
}
|
||||||
|
// есть постмортем-уведомление с анализом
|
||||||
|
var pmCount int
|
||||||
|
for _, txt := range texts {
|
||||||
|
if strings.Contains(txt, "анализ (после failed)") {
|
||||||
|
pmCount++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if pmCount != 1 {
|
||||||
|
t.Errorf("постмортем-уведомлений = %d, want ровно 1: %#v", pmCount, texts)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestWorkerTimeoutNotification — RC=-1 (таймаут dev) → уведомление о timeout.
|
// TestWorkerTimeoutNotification — RC=-1 (таймаут dev) → уведомление о timeout
|
||||||
|
// и постмортем-уведомление с анализом.
|
||||||
func TestWorkerTimeoutNotification(t *testing.T) {
|
func TestWorkerTimeoutNotification(t *testing.T) {
|
||||||
s := setupWorkerDB(t)
|
s := setupWorkerDB(t)
|
||||||
task := createReadyTask(t, s, "notif-timeout")
|
task := createReadyTask(t, s, "notif-timeout")
|
||||||
@@ -439,13 +467,21 @@ func TestWorkerTimeoutNotification(t *testing.T) {
|
|||||||
_ = w.runTask(ctx, task)
|
_ = w.runTask(ctx, task)
|
||||||
|
|
||||||
prefix := "Задача #" + strconv.FormatInt(task.ID, 10)
|
prefix := "Задача #" + strconv.FormatInt(task.ID, 10)
|
||||||
|
texts := notifTexts(n)
|
||||||
|
if len(texts) != 3 {
|
||||||
|
t.Fatalf("уведомлений = %d, want 3: %#v", len(texts), texts)
|
||||||
|
}
|
||||||
want := []string{prefix + ": running", prefix + ": timeout"}
|
want := []string{prefix + ": running", prefix + ": timeout"}
|
||||||
if got := notifTexts(n); !reflect.DeepEqual(got, want) {
|
if !reflect.DeepEqual(texts[:2], want) {
|
||||||
t.Errorf("уведомления = %#v, want %#v", got, want)
|
t.Errorf("первые уведомления = %#v, want %#v", texts[:2], want)
|
||||||
|
}
|
||||||
|
if !strings.Contains(texts[2], "🔍 анализ (после timeout)") {
|
||||||
|
t.Errorf("постмортем-уведомление = %q, want упоминание «анализ (после timeout)»", texts[2])
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestWorkerSpawnErrorNotification — сбой запуска dev → уведомление о failed.
|
// TestWorkerSpawnErrorNotification — сбой запуска dev → уведомление о failed
|
||||||
|
// и постмортем-уведомление.
|
||||||
func TestWorkerSpawnErrorNotification(t *testing.T) {
|
func TestWorkerSpawnErrorNotification(t *testing.T) {
|
||||||
s := setupWorkerDB(t)
|
s := setupWorkerDB(t)
|
||||||
task := createReadyTask(t, s, "notif-spawn")
|
task := createReadyTask(t, s, "notif-spawn")
|
||||||
@@ -463,9 +499,16 @@ func TestWorkerSpawnErrorNotification(t *testing.T) {
|
|||||||
_ = w.runTask(ctx, task)
|
_ = w.runTask(ctx, task)
|
||||||
|
|
||||||
prefix := "Задача #" + strconv.FormatInt(task.ID, 10)
|
prefix := "Задача #" + strconv.FormatInt(task.ID, 10)
|
||||||
|
texts := notifTexts(n)
|
||||||
|
if len(texts) != 3 {
|
||||||
|
t.Fatalf("уведомлений = %d, want 3: %#v", len(texts), texts)
|
||||||
|
}
|
||||||
want := []string{prefix + ": running", prefix + ": failed"}
|
want := []string{prefix + ": running", prefix + ": failed"}
|
||||||
if got := notifTexts(n); !reflect.DeepEqual(got, want) {
|
if !reflect.DeepEqual(texts[:2], want) {
|
||||||
t.Errorf("уведомления = %#v, want %#v", got, want)
|
t.Errorf("первые уведомления = %#v, want %#v", texts[:2], want)
|
||||||
|
}
|
||||||
|
if !strings.Contains(texts[2], "анализ (после failed)") {
|
||||||
|
t.Errorf("постмортем-уведомление = %q, want упоминание «анализ (после failed)»", texts[2])
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -529,9 +572,10 @@ func TestWorkerTimeout(t *testing.T) {
|
|||||||
s := setupWorkerDB(t)
|
s := setupWorkerDB(t)
|
||||||
task := createReadyTask(t, s, "slow")
|
task := createReadyTask(t, s, "slow")
|
||||||
|
|
||||||
|
runner := &mockRunnerWorker{result: &opencode.Result{RC: -1, Stdout: ""}}
|
||||||
w := &Worker{
|
w := &Worker{
|
||||||
Store: s,
|
Store: s,
|
||||||
Runner: &mockRunnerWorker{result: &opencode.Result{RC: -1, Stdout: ""}},
|
Runner: runner,
|
||||||
Worktree: t.TempDir(),
|
Worktree: t.TempDir(),
|
||||||
}
|
}
|
||||||
seedFakeRepo(t, w.Worktree, "slow")
|
seedFakeRepo(t, w.Worktree, "slow")
|
||||||
@@ -551,11 +595,20 @@ func TestWorkerTimeout(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("get traces: %v", err)
|
t.Fatalf("get traces: %v", err)
|
||||||
}
|
}
|
||||||
if len(traces) != 1 {
|
if len(traces) != 2 {
|
||||||
t.Fatalf("got %d traces, want 1", len(traces))
|
t.Fatalf("got %d traces, want 2 (dev + postmortem)", len(traces))
|
||||||
}
|
}
|
||||||
if traces[0].Status != storage.TraceTimeout {
|
if traces[0].Status != storage.TraceTimeout {
|
||||||
t.Errorf("trace status = %q, want timeout", traces[0].Status)
|
t.Errorf("trace[0] status = %q, want timeout", traces[0].Status)
|
||||||
|
}
|
||||||
|
if traces[1].Agent != postMortemAgent {
|
||||||
|
t.Errorf("trace[1] agent = %q, want postmortem", traces[1].Agent)
|
||||||
|
}
|
||||||
|
if traces[1].Status != storage.TraceSuccess {
|
||||||
|
t.Errorf("trace[1] status = %q, want success", traces[1].Status)
|
||||||
|
}
|
||||||
|
if runner.postMortemCount != 1 {
|
||||||
|
t.Errorf("postmortem запускался %d раз, want 1", runner.postMortemCount)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -563,9 +616,10 @@ func TestWorkerSpawnError(t *testing.T) {
|
|||||||
s := setupWorkerDB(t)
|
s := setupWorkerDB(t)
|
||||||
task := createReadyTask(t, s, "spawn-fail")
|
task := createReadyTask(t, s, "spawn-fail")
|
||||||
|
|
||||||
|
runner := &mockRunnerWorker{err: errors.New("opencode not found")}
|
||||||
w := &Worker{
|
w := &Worker{
|
||||||
Store: s,
|
Store: s,
|
||||||
Runner: &mockRunnerWorker{err: errors.New("opencode not found")},
|
Runner: runner,
|
||||||
Worktree: t.TempDir(),
|
Worktree: t.TempDir(),
|
||||||
}
|
}
|
||||||
seedFakeRepo(t, w.Worktree, "spawn-fail")
|
seedFakeRepo(t, w.Worktree, "spawn-fail")
|
||||||
@@ -585,11 +639,17 @@ func TestWorkerSpawnError(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("get traces: %v", err)
|
t.Fatalf("get traces: %v", err)
|
||||||
}
|
}
|
||||||
if len(traces) != 1 {
|
if len(traces) != 2 {
|
||||||
t.Fatalf("got %d traces, want 1", len(traces))
|
t.Fatalf("got %d traces, want 2 (dev-failed + postmortem)", len(traces))
|
||||||
}
|
}
|
||||||
if traces[0].Status != storage.TraceFailed {
|
if traces[0].Status != storage.TraceFailed {
|
||||||
t.Errorf("trace status = %q, want failed", traces[0].Status)
|
t.Errorf("trace[0] status = %q, want failed", traces[0].Status)
|
||||||
|
}
|
||||||
|
if traces[1].Agent != postMortemAgent {
|
||||||
|
t.Errorf("trace[1] agent = %q, want postmortem", traces[1].Agent)
|
||||||
|
}
|
||||||
|
if runner.postMortemCount != 1 {
|
||||||
|
t.Errorf("postmortem запускался %d раз, want 1", runner.postMortemCount)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -619,8 +679,8 @@ func TestWorkerNonZeroExit(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("get traces: %v", err)
|
t.Fatalf("get traces: %v", err)
|
||||||
}
|
}
|
||||||
if len(traces) != 1 {
|
if len(traces) != 2 {
|
||||||
t.Fatalf("got %d traces, want 1", len(traces))
|
t.Fatalf("got %d traces, want 2 (dev + postmortem)", len(traces))
|
||||||
}
|
}
|
||||||
if traces[0].Status != storage.TraceFailed {
|
if traces[0].Status != storage.TraceFailed {
|
||||||
t.Errorf("trace status = %q, want failed", traces[0].Status)
|
t.Errorf("trace status = %q, want failed", traces[0].Status)
|
||||||
@@ -628,6 +688,9 @@ func TestWorkerNonZeroExit(t *testing.T) {
|
|||||||
if traces[0].Output != "error" {
|
if traces[0].Output != "error" {
|
||||||
t.Errorf("output = %q, want error", traces[0].Output)
|
t.Errorf("output = %q, want error", traces[0].Output)
|
||||||
}
|
}
|
||||||
|
if traces[1].Agent != postMortemAgent {
|
||||||
|
t.Errorf("trace[1] agent = %q, want postmortem", traces[1].Agent)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestWorkerBadStatus(t *testing.T) {
|
func TestWorkerBadStatus(t *testing.T) {
|
||||||
@@ -882,3 +945,148 @@ func TestWorkerSemaphore(t *testing.T) {
|
|||||||
t.Errorf("после освобождения слота success = %d, want 2", len(success))
|
t.Errorf("после освобождения слота success = %d, want 2", len(success))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestWorkerPostMortemSkippedOnSuccess — успешный трейд НЕ запускает постмортем:
|
||||||
|
// причина анализа — только failed/timeout.
|
||||||
|
func TestWorkerPostMortemSkippedOnSuccess(t *testing.T) {
|
||||||
|
s := setupWorkerDB(t)
|
||||||
|
task := createReadyTask(t, s, "pm-ok")
|
||||||
|
|
||||||
|
runner := &mockRunnerWorker{result: &opencode.Result{RC: 0, Stdout: "done", SessionID: "s"}}
|
||||||
|
w := &Worker{
|
||||||
|
Store: s,
|
||||||
|
Runner: runner,
|
||||||
|
Worktree: t.TempDir(),
|
||||||
|
}
|
||||||
|
seedFakeRepo(t, w.Worktree, "pm-ok")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
if err := w.runTask(ctx, task); err != nil {
|
||||||
|
t.Fatalf("runTask: %v", err)
|
||||||
|
}
|
||||||
|
task, _ = s.GetTask(ctx, task.ID)
|
||||||
|
if task.Status != storage.StatusSuccess {
|
||||||
|
t.Fatalf("status = %q, want success", task.Status)
|
||||||
|
}
|
||||||
|
if runner.postMortemCount != 0 {
|
||||||
|
t.Errorf("postmortem запускался %d раз, want 0 при success", runner.postMortemCount)
|
||||||
|
}
|
||||||
|
traces, _ := s.GetTraces(ctx, task.ID)
|
||||||
|
for _, tr := range traces {
|
||||||
|
if tr.Agent == postMortemAgent {
|
||||||
|
t.Errorf("есть неожиданный postmortem-trace при success")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestWorkerPostMortemFailureDoesNotChangeStatus — сбой самого постмортема не
|
||||||
|
// влияет на статус задачи (остаётся failed) и фиксируется как failed-трасса.
|
||||||
|
func TestWorkerPostMortemFailureDoesNotChangeStatus(t *testing.T) {
|
||||||
|
s := setupWorkerDB(t)
|
||||||
|
task := createReadyTask(t, s, "pm-fail")
|
||||||
|
|
||||||
|
runner := &mockRunnerWorker{
|
||||||
|
// dev падает при спавне → failed; постмортем тоже падает.
|
||||||
|
err: errors.New("opencode not found"),
|
||||||
|
postMortemResult: &opencode.Result{RC: 1, Stdout: ""},
|
||||||
|
}
|
||||||
|
w := &Worker{
|
||||||
|
Store: s,
|
||||||
|
Runner: runner,
|
||||||
|
Worktree: t.TempDir(),
|
||||||
|
}
|
||||||
|
seedFakeRepo(t, w.Worktree, "pm-fail")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
_ = w.runTask(ctx, task)
|
||||||
|
|
||||||
|
task, _ = s.GetTask(ctx, task.ID)
|
||||||
|
if task.Status != storage.StatusFailed {
|
||||||
|
t.Errorf("status = %q, want failed (постмортем не должен менять статус)", task.Status)
|
||||||
|
}
|
||||||
|
traces, err := s.GetTraces(ctx, task.ID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get traces: %v", err)
|
||||||
|
}
|
||||||
|
if len(traces) != 2 {
|
||||||
|
t.Fatalf("traces = %d, want 2 (dev-failed + postmortem)", len(traces))
|
||||||
|
}
|
||||||
|
if traces[1].Agent != postMortemAgent {
|
||||||
|
t.Errorf("trace[1] agent = %q, want postmortem", traces[1].Agent)
|
||||||
|
}
|
||||||
|
if traces[1].Status != storage.TraceFailed {
|
||||||
|
t.Errorf("trace[1] status = %q, want failed (сбой постмортема)", traces[1].Status)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestWorkerPostMortemNotRepeated — если у задачи уже есть postmortem-trace
|
||||||
|
// (например, от прошлого прогона), повторный постмортем не запускается.
|
||||||
|
func TestWorkerPostMortemNotRepeated(t *testing.T) {
|
||||||
|
s := setupWorkerDB(t)
|
||||||
|
task := createReadyTask(t, s, "pm-repeat")
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
if _, err := s.AppendTrace(ctx, &storage.Trace{
|
||||||
|
TaskID: task.ID,
|
||||||
|
Agent: postMortemAgent,
|
||||||
|
Prompt: "старый анализ",
|
||||||
|
Output: "старое резюме",
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("seed postmortem trace: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
runner := &mockRunnerWorker{err: errors.New("opencode not found")}
|
||||||
|
w := &Worker{
|
||||||
|
Store: s,
|
||||||
|
Runner: runner,
|
||||||
|
Worktree: t.TempDir(),
|
||||||
|
}
|
||||||
|
seedFakeRepo(t, w.Worktree, "pm-repeat")
|
||||||
|
|
||||||
|
_ = w.runTask(ctx, task)
|
||||||
|
|
||||||
|
task, _ = s.GetTask(ctx, task.ID)
|
||||||
|
if task.Status != storage.StatusFailed {
|
||||||
|
t.Fatalf("status = %q, want failed", task.Status)
|
||||||
|
}
|
||||||
|
if runner.postMortemCount != 0 {
|
||||||
|
t.Errorf("postmortem запускался %d раз, want 0 (уже был trace)", runner.postMortemCount)
|
||||||
|
}
|
||||||
|
// количество postmortem-трасс не выросло
|
||||||
|
traces, _ := s.GetTraces(ctx, task.ID)
|
||||||
|
var pm int
|
||||||
|
for _, tr := range traces {
|
||||||
|
if tr.Agent == postMortemAgent {
|
||||||
|
pm++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if pm != 1 {
|
||||||
|
t.Errorf("postmortem-трасс = %d, want 1 (без дубля)", pm)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestRenderPostMortemPrompt — промпт постмортема включает задачу и сессии.
|
||||||
|
func TestRenderPostMortemPrompt(t *testing.T) {
|
||||||
|
tr := storage.Trace{
|
||||||
|
Agent: "dev",
|
||||||
|
Status: storage.TraceTimeout,
|
||||||
|
Prompt: "промпт dev",
|
||||||
|
Output: "вывод dev",
|
||||||
|
}
|
||||||
|
prompt, err := RenderPostMortemPrompt(PostMortemPromptData{
|
||||||
|
Title: "Таймаут-задача",
|
||||||
|
Goal: "сделать",
|
||||||
|
Repos: []string{"calc"},
|
||||||
|
AC: "работает",
|
||||||
|
Status: storage.StatusTimeout,
|
||||||
|
Sessions: postMortemsText([]storage.Trace{tr}),
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("render: %v", err)
|
||||||
|
}
|
||||||
|
for _, want := range []string{"Таймаут-задача", "timeout", "=== Агент: dev", "промпт dev", "вывод dev"} {
|
||||||
|
if !strings.Contains(prompt, want) {
|
||||||
|
t.Errorf("промпт не содержит %q", want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user