2 Commits

Author SHA1 Message Date
ki.sagidullin
6c0b903794 feat(analyst): этапы задачи — аналитик раскладывает на steps с критериями готовности, dev идёт по ним
Some checks failed
CI / test (push) Failing after 1m52s
CI / build-and-package (amd64, linux) (push) Failing after 1m3s
CI / build-and-package (amd64, windows) (push) Successful in 30s
2026-08-24 18:05:06 +05:00
ki.sagidullin
6546f18348 fix(worker): постмортем получает транскрипт сессий dev/reviewer по session_id
Some checks failed
CI / test (push) Failing after 1m52s
CI / build-and-package (amd64, linux) (push) Failing after 1m1s
CI / build-and-package (amd64, windows) (push) Successful in 29s
Раньше постмортем-агенту передавался только session_id, без содержимого
сессии — анализировать было нечего. Теперь Runner умеет читать сообщения
сессии (Messages → renderTranscript) через SessionMessages, а postMortem
подставляет реальные шаги агента (сообщения/tool-вызовы) в промпт.
Пустая/недоступная сессия помечается явно. Сбои чтения не фатальны.
2026-08-24 16:48:43 +05:00
13 changed files with 507 additions and 66 deletions

View File

@@ -57,6 +57,7 @@ type AnalystResponse struct {
Repos json.RawMessage `json:"repos"` // список репо (основной); устойчив к строке
Why string `json:"why"`
AC string `json:"ac"`
Steps []storage.Step `json:"steps"`
Questions []string `json:"questions"`
ChatReply string `json:"chat_reply"`
AbortReason string `json:"abort_reason"`
@@ -86,6 +87,22 @@ func (r *AnalystResponse) reposList() []string {
return out
}
// stepsList возвращает этапы: модель может вернуть массив объектов либо
// пустой/отсутствующий — тогда nil. Записи без title отбрасываются.
func (r *AnalystResponse) stepsList() []storage.Step {
if r.Steps == nil {
return nil
}
var out []storage.Step
for _, st := range r.Steps {
if st.Title == "" {
continue
}
out = append(out, st)
}
return out
}
// Decide реализует core.Decider через открытый код.
func (a *Analyst) Decide(ctx context.Context, history []core.Message, draft storage.Task, force bool) (core.Decision, error) {
agent := a.Agent
@@ -109,6 +126,7 @@ func (a *Analyst) Decide(ctx context.Context, history []core.Message, draft stor
Repos: draft.EffectiveRepos(),
Why: draft.Why,
AC: draft.AC,
Steps: draft.Steps,
History: hist,
Force: force,
}
@@ -176,6 +194,9 @@ func (a *Analyst) Decide(ctx context.Context, history []core.Message, draft stor
if ar.AC != "" {
dec.Draft.AC = ar.AC
}
if steps := ar.stepsList(); len(steps) > 0 {
dec.Draft.Steps = steps
}
return dec, nil
}
@@ -204,7 +225,7 @@ func validateResponse(ar *AnalystResponse) error {
return fmt.Errorf("phase=ask, но нет ни chat_reply, ни questions")
}
case "propose":
if ar.Title == "" && ar.Goal == "" && len(ar.Repos) == 0 && ar.Why == "" && ar.AC == "" {
if ar.Title == "" && ar.Goal == "" && len(ar.Repos) == 0 && ar.Why == "" && ar.AC == "" && len(ar.stepsList()) == 0 {
return fmt.Errorf("phase=propose, но нет ни одного изменённого поля")
}
case "ready":
@@ -257,6 +278,19 @@ func formatVerdict(ar *AnalystResponse) string {
b.WriteString(", ac=")
b.WriteString(ar.AC)
}
if steps := ar.stepsList(); len(steps) > 0 {
b.WriteString(", steps=[")
var parts []string
for _, st := range steps {
s := st.Title
if st.AC != "" {
s += " → " + st.AC
}
parts = append(parts, s)
}
b.WriteString(strings.Join(parts, " | "))
b.WriteString("]")
}
if ar.AbortReason != "" {
b.WriteString(", abort_reason=")
b.WriteString(ar.AbortReason)

View File

@@ -240,6 +240,52 @@ func TestDecideProposeStringRepos(t *testing.T) {
}
}
// TestDecideProposeSteps — аналитик разложил задачу на этапы с критериями.
func TestDecideProposeSteps(t *testing.T) {
a := &Analyst{Runner: &mockRunner{result: &opencode.Result{
RC: 0,
Stdout: `{"type":"text","part":{"text":"{\"phase\":\"propose\",\"title\":\"Калькулятор\",\"steps\":[{\"title\":\"Модель\",\"ac\":\"операции + - * /\"},{\"title\":\"UI\"}]}"}}`,
}}, Worktree: "/tmp"}
history := []core.Message{{Role: "user", Content: "Сделай калькулятор"}}
dec, err := a.Decide(context.Background(), history, storage.Task{}, false)
if err != nil {
t.Fatalf("Decide err: %v", err)
}
if dec.Phase != "propose" {
t.Errorf("Phase = %q, want propose", dec.Phase)
}
if len(dec.Draft.Steps) != 2 {
t.Fatalf("len(Steps) = %d, want 2", len(dec.Draft.Steps))
}
if dec.Draft.Steps[0].Title != "Модель" || dec.Draft.Steps[0].AC != "операции + - * /" {
t.Errorf("Steps[0] = %q/%q, want Модель/операции + - * /", dec.Draft.Steps[0].Title, dec.Draft.Steps[0].AC)
}
if dec.Draft.Steps[1].Title != "UI" || dec.Draft.Steps[1].AC != "" {
t.Errorf("Steps[1] = %q/%q, want UI/(пусто)", dec.Draft.Steps[1].Title, dec.Draft.Steps[1].AC)
}
}
// TestProposeOnlyStepsValid — propose меняет только steps → валидно.
func TestProposeOnlyStepsValid(t *testing.T) {
a := &Analyst{Runner: &mockRunner{result: &opencode.Result{
RC: 0,
Stdout: `{"type":"text","part":{"text":"{\"phase\":\"propose\",\"steps\":[{\"title\":\"Шаг 1\"}],\"chat_reply\":\"Разбил на этапы\"}"}}`,
}}, Worktree: "/tmp"}
history := []core.Message{{Role: "user", Content: "test"}}
dec, err := a.Decide(context.Background(), history, storage.Task{}, false)
if err != nil {
t.Fatalf("Decide err: %v", err)
}
if dec.Phase != "propose" {
t.Errorf("Phase = %q, want propose", dec.Phase)
}
if len(dec.Draft.Steps) != 1 {
t.Fatalf("len(Steps) = %d, want 1", len(dec.Draft.Steps))
}
}
// TestFormatVerdict — человекочитаемое описание вердикта аналитика.
func TestFormatVerdict(t *testing.T) {
tests := []struct {

View File

@@ -3,6 +3,8 @@ package analyst
import (
"strings"
"text/template"
"github.com/kamelion/ratatoskr-go/internal/storage"
)
// promptTemplate — шаблон промпта для аналитика (opencode analyst-agent).
@@ -23,6 +25,7 @@ var promptTemplate = template.Must(template.New("analyst").Parse(`Ты — ан
| repos | список репозиториев (имена на git-хосте; для связанных — все сразу) |
| why | зачем это нужно, контекст |
| ac | acceptance criteria — конкретный результат, что считается готовым |
| steps | (опционально) разбиение задачи на этапы: список {title, ac} с критерием готовности каждого этапа |
Изменяй в JSON только те поля, которые надо поменять; что менять не надо — пустой строкой.
@@ -51,6 +54,12 @@ var promptTemplate = template.Must(template.New("analyst").Parse(`Ты — ан
{{- else}} repos: (не задано){{end}}
{{if .Why}} why: {{.Why}}{{else}} why: (не задано){{end}}
{{if .AC}} ac: {{.AC}}{{else}} ac: (не задано){{end}}
{{if .Steps}}
steps:
{{- range .Steps}}
- {{.Title}}{{if .AC}}{{.AC}}{{end}}
{{- end}}
{{- else}} steps: (не задано){{end}}
**Ответь строго JSON-объектом, без лишнего текста:**
{
@@ -60,6 +69,7 @@ var promptTemplate = template.Must(template.New("analyst").Parse(`Ты — ан
"repos": ["имя_репо_1", "имя_репо_2"],
"why": "зачем (только если меняешь)",
"ac": "критерии (только если меняешь)",
"steps": [{"title": "этап 1", "ac": "критерий этапа 1"}],
"questions": ["вопрос 1", "вопрос 2"],
"chat_reply": "твой ответ пользователю (на русском, естественно)",
"abort_reason": "если phase=abort — причина"
@@ -73,6 +83,7 @@ type TemplateData struct {
Repos []string
Why string
AC string
Steps []storage.Step
History string // отформатированная переписка
Force bool
}

View File

@@ -7,6 +7,7 @@ import (
"fmt"
"io"
"net/http"
"strings"
"time"
)
@@ -349,6 +350,59 @@ func assistantVerdict(msgs []v2Message, since int64) (texts []string, usedReason
return reasoning, len(reasoning) > 0
}
// transcriptRole возвращает человекочитаемую подпись роли сообщения сессии.
func transcriptRole(t string) string {
switch t {
case "user":
return "Пользователь"
case "assistant":
return "Ассистент"
case "tool":
return "Инструмент"
case "system":
return "Система"
default:
return t
}
}
// finishLabel — подпись для finish-reason (пусто/: опускаем).
func finishLabel(f string) string {
if f == "" || f == "stop" {
return ""
}
return f
}
// renderTranscript форматирует сообщения сессии (API отдаёт новыми первыми)
// в хронологическом порядке: роль, текст (и tool-вызовы), с обрезкой.
func renderTranscript(msgs []v2Message) string {
var b strings.Builder
for i := len(msgs) - 1; i >= 0; i-- {
m := &msgs[i]
b.WriteString("\n== " + transcriptRole(m.Type))
if m.Type == "assistant" {
if m.Model != nil {
b.WriteString(" (" + m.Model.String() + ")")
}
if fl := finishLabel(m.Finish); fl != "" {
b.WriteString(" finish=" + fl)
}
}
b.WriteString(" ==\n")
for _, p := range m.Content {
if p.Text == "" {
continue
}
b.WriteString(p.Text)
if !strings.HasSuffix(p.Text, "\n") {
b.WriteString("\n")
}
}
}
return b.String()
}
func truncateStr(s string, n int) string {
if len(s) <= n {
return s

View File

@@ -90,6 +90,30 @@ func (r *Runner) Run(ctx context.Context, prompt, cwd, agent, sessionID string)
return r.awaitVerdict(ctx, c, sid, agent, prompt)
}
// SessionMessages возвращает транскрипт сессии (шаги агента: сообщения,
// тексты, tool-вызовы) в хронологическом порядке для постмортем-анализа.
// sessionID пустой → пустая строка без ошибки. Ошибки чтения возвращаются
// как есть — вызывающий (постмортем) решает, логировать и продолжить.
func (r *Runner) SessionMessages(ctx context.Context, cwd, sessionID string) (string, error) {
r.defaults()
if r.Pool == nil {
return "", fmt.Errorf("opencode: Pool не задан (API-режим обязателен)")
}
if sessionID == "" {
return "", nil
}
srv, err := r.Pool.Ensure(ctx, cwd)
if err != nil {
return "", err
}
c := &Client{BaseURL: srv.Addr(), Password: srv.Password, Directory: srv.Dir, Debug: r.Debug}
msgs, err := c.Messages(ctx, sessionID)
if err != nil {
return "", err
}
return renderTranscript(msgs), nil
}
// settlePolls — сколько подряд опросов должно подтвердить завершение ответа,
// прежде чем считать вердикт финальным (устойчивость к гонке между удалением
// сессии из активных дренажей и финализацией последнего сообщения).

View File

@@ -40,6 +40,12 @@ func IsTerminal(s Status) bool {
return model.IsTerminal(s)
}
// Step — этап задачи с собственным критерием готовности.
type Step struct {
Title string `json:"title"`
AC string `json:"ac"` // acceptance criterion этапа
}
// Task — запись задачи в БД.
type Task struct {
ID int64 `json:"id"`
@@ -50,6 +56,7 @@ type Task struct {
Repos []string `json:"repos"` // список репозиториев (основной)
Why string `json:"why"`
AC string `json:"ac"` // acceptance criteria
Steps []Step `json:"steps"` // этапы задачи (опционально)
TaskTag string `json:"task_tag"` // UUID, стабильный на всю жизнь
Status Status `json:"status"`
CreatedAt SQLiteTime `json:"created_at"`
@@ -88,6 +95,25 @@ func (t *Task) SetReposFromDB(repos string) {
_ = json.Unmarshal([]byte(repos), &t.Repos)
}
// StepsJoined возвращает steps как одну строку (JSON-массив) для хранения в БД.
// Пустой список → пустая строка.
func (t *Task) StepsJoined() string {
if len(t.Steps) == 0 {
return ""
}
b, _ := json.Marshal(t.Steps)
return string(b)
}
// SetStepsFromDB заполняет Steps из сохранённой строки (JSON).
func (t *Task) SetStepsFromDB(steps string) {
if steps == "" {
t.Steps = nil
return
}
_ = json.Unmarshal([]byte(steps), &t.Steps)
}
// TraceStatus — алиас доменного статуса трассировки.
type TraceStatus = model.TraceStatus

View File

@@ -173,6 +173,13 @@ func (s *Storage) migrate(ctx context.Context) error {
return fmt.Errorf("%w: migrate add repos: %w", ErrDB, err)
}
}
// Доп. колонка steps (этапы задачи). Idempotent.
if _, err := s.db.ExecContext(ctx,
`ALTER TABLE tasks ADD COLUMN steps TEXT NOT NULL DEFAULT ''`); err != nil {
if !isDuplicateColumn(err) {
return fmt.Errorf("%w: migrate add steps: %w", ErrDB, err)
}
}
return nil
}

View File

@@ -150,6 +150,44 @@ func TestUpdateTaskNotFound(t *testing.T) {
}
}
func TestCreateAndGetTask_Steps(t *testing.T) {
s, ctx := setupTestDB(t)
task := &Task{
ChatID: "tg://steps",
Title: "Steps task",
TaskTag: "steps-1",
Steps: []Step{
{Title: "Реализовать модель", AC: "структура готова"},
{Title: "Добавить API", AC: "эндпоинт отвечает"},
},
}
id, err := s.CreateTask(ctx, task)
if err != nil {
t.Fatalf("CreateTask: %v", err)
}
got, err := s.GetTask(ctx, id)
if err != nil {
t.Fatalf("GetTask: %v", err)
}
if len(got.Steps) != 2 {
t.Fatalf("steps len = %d, want 2", len(got.Steps))
}
if got.Steps[0].Title != "Реализовать модель" || got.Steps[0].AC != "структура готова" {
t.Fatalf("steps[0] = %q / %q, want модель / структура готова", got.Steps[0].Title, got.Steps[0].AC)
}
// апдейт этапов
got.Steps = append(got.Steps, Step{Title: "Ревью", AC: "пройден review"})
if err := s.UpdateTask(ctx, got); err != nil {
t.Fatalf("UpdateTask steps: %v", err)
}
got2, _ := s.GetTask(ctx, id)
if len(got2.Steps) != 3 {
t.Fatalf("steps len after update = %d, want 3", len(got2.Steps))
}
}
func TestListTasks(t *testing.T) {
s, ctx := setupTestDB(t)
for i := 0; i < 5; i++ {

View File

@@ -30,10 +30,10 @@ func (s *Storage) CreateTask(ctx context.Context, t *Task) (int64, error) {
}
now := Now()
res, err := s.db.ExecContext(ctx, `
INSERT INTO tasks (chat_id, title, goal, repo, repos, why, ac, task_tag, status, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
INSERT INTO tasks (chat_id, title, goal, repo, repos, why, ac, steps, task_tag, status, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
t.ChatID, t.Title, t.Goal, t.Repo, t.ReposJoined(), t.Why, t.AC,
t.TaskTag, StatusDraft, now, now,
t.StepsJoined(), t.TaskTag, StatusDraft, now, now,
)
if err != nil {
return 0, fmt.Errorf("%w: create task: %w", ErrDB, err)
@@ -53,11 +53,12 @@ func (s *Storage) CreateTask(ctx context.Context, t *Task) (int64, error) {
func (s *Storage) GetTask(ctx context.Context, id int64) (*Task, error) {
t := &Task{}
var reposStr string
var stepsStr string
err := s.db.QueryRowContext(ctx, `
SELECT id, chat_id, title, goal, repo, repos, why, ac, task_tag, status, created_at, updated_at
SELECT id, chat_id, title, goal, repo, repos, why, ac, steps, task_tag, status, created_at, updated_at
FROM tasks WHERE id = ?`, id).Scan(
&t.ID, &t.ChatID, &t.Title, &t.Goal, &t.Repo, &reposStr,
&t.Why, &t.AC, &t.TaskTag, &t.Status, &t.CreatedAt, &t.UpdatedAt,
&t.Why, &t.AC, &stepsStr, &t.TaskTag, &t.Status, &t.CreatedAt, &t.UpdatedAt,
)
if err == sql.ErrNoRows {
return nil, fmt.Errorf("%w: task %d", ErrNotFound, id)
@@ -66,6 +67,7 @@ func (s *Storage) GetTask(ctx context.Context, id int64) (*Task, error) {
return nil, fmt.Errorf("%w: get task %d: %w", ErrDB, id, err)
}
t.SetReposFromDB(reposStr)
t.SetStepsFromDB(stepsStr)
return t, nil
}
@@ -94,9 +96,9 @@ func (s *Storage) UpdateTask(ctx context.Context, t *Task) error {
now := Now()
res, err := s.db.ExecContext(ctx, `
UPDATE tasks
SET title=?, goal=?, repo=?, repos=?, why=?, ac=?, status=?, updated_at=?
SET title=?, goal=?, repo=?, repos=?, why=?, ac=?, steps=?, status=?, updated_at=?
WHERE id=?`,
t.Title, t.Goal, t.Repo, t.ReposJoined(), t.Why, t.AC, t.Status, now, t.ID,
t.Title, t.Goal, t.Repo, t.ReposJoined(), t.Why, t.AC, t.StepsJoined(), t.Status, now, t.ID,
)
if err != nil {
return fmt.Errorf("%w: update task %d: %w", ErrDB, t.ID, err)
@@ -114,13 +116,14 @@ func (s *Storage) UpdateTask(ctx context.Context, t *Task) error {
func (s *Storage) GetActiveTaskByChatID(ctx context.Context, chatID string) (*Task, error) {
t := &Task{}
var reposStr string
var stepsStr string
err := s.db.QueryRowContext(ctx, `
SELECT id, chat_id, title, goal, repo, repos, why, ac, task_tag, status, created_at, updated_at
SELECT id, chat_id, title, goal, repo, repos, why, ac, steps, task_tag, status, created_at, updated_at
FROM tasks
WHERE chat_id = ? AND status NOT IN ('success','cancelled','aborted','closed')
ORDER BY updated_at DESC LIMIT 1`, chatID).Scan(
&t.ID, &t.ChatID, &t.Title, &t.Goal, &t.Repo, &reposStr,
&t.Why, &t.AC, &t.TaskTag, &t.Status, &t.CreatedAt, &t.UpdatedAt)
&t.Why, &t.AC, &stepsStr, &t.TaskTag, &t.Status, &t.CreatedAt, &t.UpdatedAt)
if err == sql.ErrNoRows {
return nil, fmt.Errorf("%w: no active task for chat %s", ErrNotFound, chatID)
}
@@ -128,6 +131,7 @@ func (s *Storage) GetActiveTaskByChatID(ctx context.Context, chatID string) (*Ta
return nil, fmt.Errorf("%w: get active task %s: %w", ErrDB, chatID, err)
}
t.SetReposFromDB(reposStr)
t.SetStepsFromDB(stepsStr)
return t, nil
}
@@ -149,7 +153,7 @@ func (s *Storage) ListTasks(ctx context.Context, filter TaskFilter) ([]*Task, er
args = append(args, filter.Limit, filter.Offset)
rows, err := s.db.QueryContext(ctx, `
SELECT id, chat_id, title, goal, repo, repos, why, ac, task_tag, status, created_at, updated_at
SELECT id, chat_id, title, goal, repo, repos, why, ac, steps, task_tag, status, created_at, updated_at
FROM tasks WHERE `+where+` ORDER BY updated_at DESC LIMIT ? OFFSET ?`, args...)
if err != nil {
return nil, fmt.Errorf("%w: list tasks: %w", ErrDB, err)
@@ -160,11 +164,13 @@ func (s *Storage) ListTasks(ctx context.Context, filter TaskFilter) ([]*Task, er
for rows.Next() {
t := &Task{}
var reposStr string
var stepsStr string
if err := rows.Scan(&t.ID, &t.ChatID, &t.Title, &t.Goal, &t.Repo, &reposStr,
&t.Why, &t.AC, &t.TaskTag, &t.Status, &t.CreatedAt, &t.UpdatedAt); err != nil {
&t.Why, &t.AC, &stepsStr, &t.TaskTag, &t.Status, &t.CreatedAt, &t.UpdatedAt); err != nil {
return nil, fmt.Errorf("%w: scan task: %w", ErrDB, err)
}
t.SetReposFromDB(reposStr)
t.SetStepsFromDB(stepsStr)
tasks = append(tasks, t)
}
return tasks, rows.Err()

View File

@@ -18,6 +18,9 @@ const postMortemAgent = "postmortem"
// не превращался в полные транскрипты и не переполнял контекст модели).
const traceOutputMax = 6000
// sessionTranscriptMax — обрезка транскрипта сессии dev/reviewer в промпте.
const sessionTranscriptMax = 20000
// postMortemPromptTemplate — промпт для постмортем-агента после failed/timeout:
// задача + сессии dev/reviewer. Ожидается резюме простым текстом на русском.
var postMortemPromptTemplate = template.Must(template.New("postmortem").Parse(`Ты — постмортем-аналитик в конвейере Ratatoskr. Задача завершилась неудачей ({{.Status}}). Проанализируй сессии агентов dev/reviewer и дай резюме: почему так случилось и что сделать, чтобы не повторялось.
@@ -70,8 +73,11 @@ func RenderPostMortemPrompt(data PostMortemPromptData) (string, error) {
return buf.String(), nil
}
// postMortemsText форматирует сессии dev/reviewer в секцию промпта.
func postMortemsText(traces []storage.Trace) string {
// postMortemsText форматирует сессии dev/reviewer в секцию промпта: поля
// трассы (session_id, промпт, вывод) + транскрипт сессии (шаги агента).
// transcripts — sessionID → транскрипт (пустая строка — в сессии нет сообщений);
// отсутствие ключа — транскрипт недоступен (сбой чтения).
func postMortemsText(traces []storage.Trace, transcripts map[string]string) string {
var b strings.Builder
for _, tr := range traces {
b.WriteString("\n=== Агент: " + tr.Agent + " (статус " + string(tr.Status) + ") ===\n")
@@ -88,6 +94,19 @@ func postMortemsText(traces []storage.Trace) string {
b.WriteString(truncateTrace(tr.Output, traceOutputMax))
b.WriteString("\n")
}
if tr.SessionID != "" {
tx, ok := transcripts[tr.SessionID]
b.WriteString("-- Транскрипт сессии (шаги агента) --\n")
switch {
case !ok:
b.WriteString("(транскрипт сессии недоступен)\n")
case strings.TrimSpace(tx) == "":
b.WriteString("(в сессии нет сообщений — агент не сделал ни одного шага)\n")
default:
b.WriteString(truncateTrace(tx, sessionTranscriptMax))
b.WriteString("\n")
}
}
}
if b.Len() == 0 {
return "(сессии dev/reviewer не найдены — вероятна инфраструктурная ошибка до запуска агентов)"
@@ -144,6 +163,23 @@ func (w *Worker) postMortem(ctx context.Context, task *storage.Task) {
}
}
// Достаём транскрипты сессий dev/reviewer (по session_id) — реальные шаги
// агента (сообщения/tool-вызовы), которые не попадают в трассу (там только
// финальный вывод). Сбои чтения не фатальны: промпт соберётся без них.
transcripts := make(map[string]string, len(sessions))
for i := range sessions {
tr := &sessions[i]
if tr.SessionID == "" {
continue
}
tx, tErr := w.Runner.SessionMessages(w.runCtx(ctx, task.ID), w.Worktree, tr.SessionID)
if tErr != nil {
log.Printf("worker: task %d: постмортем: транскрипт %s (%s): %v", task.ID, tr.Agent, tr.SessionID, tErr)
continue
}
transcripts[tr.SessionID] = tx
}
prompt, pErr := RenderPostMortemPrompt(PostMortemPromptData{
Title: task.Title,
Goal: task.Goal,
@@ -151,7 +187,7 @@ func (w *Worker) postMortem(ctx context.Context, task *storage.Task) {
Why: task.Why,
AC: task.AC,
Status: task.Status,
Sessions: postMortemsText(sessions),
Sessions: postMortemsText(sessions, transcripts),
})
if pErr != nil {
log.Printf("worker: task %d: постмортем: рендер промпта: %v", task.ID, pErr)

View File

@@ -3,6 +3,8 @@ package worker
import (
"strings"
"text/template"
"github.com/kamelion/ratatoskr-go/internal/storage"
)
// devPromptTemplate — промпт для dev-агента при запуске задачи.
@@ -20,6 +22,12 @@ var devPromptTemplate = template.Must(template.New("dev").Parse(`Ты — dev-а
{{if .Why}}Зачем: {{.Why}}{{end}}
{{if .AC}}Критерии готовности:
{{.AC}}{{end}}
{{if .Steps}}
**Этапы (выполняй по порядку, у каждого свой критерий готовности):**
{{- range .Steps}}
{{.Title}}{{if .AC}} — готово, когда: {{.AC}}{{end}}
{{- end}}
{{end}}
**Инструкции:**
1. Рабочий каталог — общий корень, в котором лежат все репозитории по именам.
@@ -38,6 +46,7 @@ type DevPromptData struct {
Repos []string
Why string
AC string
Steps []storage.Step
// ReviewFeedback — замечания ревьюера при повторном прогоне dev
// (не пусто → dev должен исправить именно это).

View File

@@ -18,6 +18,9 @@ import (
// OpenCodeRunner — интерфейс для opencode (подменяемый в тестах).
type OpenCodeRunner interface {
Run(ctx context.Context, prompt, cwd, agent, sessionID string) (*opencode.Result, error)
// SessionMessages возвращает транскрипт сессии (шаги агента) по sessionID
// для постмортем-анализа. cwd — каталог, где живёт сервер пула сессии.
SessionMessages(ctx context.Context, cwd, sessionID string) (string, error)
}
// PollTaskFunc — callback для обработки готовой задачи (подменяемый в тестах).
@@ -252,6 +255,7 @@ func (w *Worker) runTask(ctx context.Context, task *storage.Task) (err error) {
Repos: repos,
Why: task.Why,
AC: task.AC,
Steps: task.Steps,
Branch: branch,
ReviewFeedback: reviewFeedbackList(branch, feedback),
}

View File

@@ -59,6 +59,11 @@ type mockRunnerWorker struct {
// резюме. postMortemCount — сколько раз постмортем вызывался.
postMortemResult *opencode.Result
postMortemCount int
// sessionTranscripts — sessionID → транскрипт для постмортем-агента;
// sessionMsgsErr — ошибка чтения транскрипта.
sessionTranscripts map[string]string
sessionMsgsErr error
}
func (m *mockRunnerWorker) Run(_ context.Context, _, _, agent, _ string) (*opencode.Result, error) {
@@ -86,6 +91,13 @@ func (m *mockRunnerWorker) Run(_ context.Context, _, _, agent, _ string) (*openc
return m.result, m.err
}
func (m *mockRunnerWorker) SessionMessages(_ context.Context, _, sessionID string) (string, error) {
if m.sessionMsgsErr != nil {
return "", m.sessionMsgsErr
}
return m.sessionTranscripts[sessionID], nil
}
// reviewFailedRunner возвращает вердикт not-passed с комментариями.
func reviewFailedRunner() *opencode.Result {
return &opencode.Result{RC: 0, Stdout: `{"passed":false,"critical_issues":[],"solid_violations":["DIP: высокая связанность"],"comments":["исправь связанность"]}`}
@@ -713,6 +725,7 @@ func TestWorkerBadStatus(t *testing.T) {
}
}
// TestWorkerPromptRendered — dev-промпт собирается из полей задачи.
func TestWorkerPromptRendered(t *testing.T) {
s := setupWorkerDB(t)
task := createReadyTask(t, s, "prompt-test")
@@ -746,6 +759,63 @@ func TestWorkerPromptRendered(t *testing.T) {
}
}
// TestWorkerPromptIncludesSteps — этапы задачи попадают в dev-промпт по порядку
// с критериями готовности.
func TestWorkerPromptIncludesSteps(t *testing.T) {
s := setupWorkerDB(t)
task := createReadyTask(t, s, "steps-test")
ctx := context.Background()
task.Steps = []storage.Step{
{Title: "Модель", AC: "операции готовы"},
{Title: "UI"},
}
if err := s.UpdateTask(ctx, task); err != nil {
t.Fatalf("update steps: %v", err)
}
w := &Worker{
Store: s,
Runner: &mockRunnerWorker{result: &opencode.Result{RC: 0, Stdout: "ok", SessionID: "s"}},
Worktree: t.TempDir(),
}
seedFakeRepo(t, w.Worktree, "steps-test")
_ = w.runTask(ctx, task)
traces, err := s.GetTraces(ctx, task.ID)
if err != nil {
t.Fatalf("get traces: %v", err)
}
tr := traces[0]
for _, want := range []string{"Этапы", "Модель", "операции готовы", "UI"} {
if !strings.Contains(tr.Prompt, want) {
t.Errorf("dev-промпт не содержит %q", want)
}
}
}
// TestRenderDevPromptSteps — прямой рендер dev-промпта с этапами.
func TestRenderDevPromptSteps(t *testing.T) {
prompt, err := RenderDevPrompt(DevPromptData{
Title: "Калькулятор",
AC: "работает",
Steps: []storage.Step{
{Title: "Модель", AC: "операции готовы"},
{Title: "UI"},
},
Branch: "feat/abc",
})
if err != nil {
t.Fatalf("render: %v", err)
}
for _, want := range []string{"Этапы", "Модель — готово, когда: операции готовы", "UI"} {
if !strings.Contains(prompt, want) {
t.Errorf("prompt не содержит %q", want)
}
}
}
func TestFeatureBranchName(t *testing.T) {
cases := []struct {
tag string
@@ -1065,32 +1135,108 @@ func TestWorkerPostMortemNotRepeated(t *testing.T) {
}
}
// TestRenderPostMortemPrompt — промпт постмортема включает задачу и сессии.
// TestWorkerPostMortemUsesTranscript — постмортем достаёт транскрипт dev-сессии
// (по session_id) и вставляет шаги агента в промпт постмортем-агента.
func TestWorkerPostMortemUsesTranscript(t *testing.T) {
s := setupWorkerDB(t)
task := createReadyTask(t, s, "pm-transcript")
runner := &mockRunnerWorker{
result: &opencode.Result{RC: -1, Stdout: "", SessionID: "sess-dev"},
sessionTranscripts: map[string]string{
"sess-dev": "[Инструмент]\nчитает requirements.md\n[Ассистент]\nправлю main.go",
},
}
w := &Worker{
Store: s,
Runner: runner,
Agent: "dev",
Worktree: t.TempDir(),
}
seedFakeRepo(t, w.Worktree, "pm-transcript")
ctx := context.Background()
_ = w.runTask(ctx, task)
task, _ = s.GetTask(ctx, task.ID)
if task.Status != storage.StatusTimeout {
t.Fatalf("status = %q, want timeout", task.Status)
}
traces, err := s.GetTraces(ctx, task.ID)
if err != nil {
t.Fatalf("get traces: %v", err)
}
var pm *storage.Trace
for _, tr := range traces {
if tr.Agent == postMortemAgent {
pm = tr
}
}
if pm == nil {
t.Fatal("нет постмортем-трассы")
}
for _, want := range []string{"Транскрипт сессии", "читает requirements.md", "правлю main.go"} {
if !strings.Contains(pm.Prompt, want) {
t.Errorf("промпт постмортема не содержит %q", want)
}
}
}
// TestRenderPostMortemPrompt — промпт постмортема включает задачу, сессии и
// транскрипт сессий (шаги агентов), а не только финальный вывод.
func TestRenderPostMortemPrompt(t *testing.T) {
tr := storage.Trace{
Agent: "dev",
Status: storage.TraceTimeout,
SessionID: "sess-dev-1",
Prompt: "промпт dev",
Output: "вывод dev",
}
transcripts := map[string]string{
"sess-dev-1": "[Инструмент]\nпрочитал файл a.go\n[Ассистент]\nправлю код",
}
prompt, err := RenderPostMortemPrompt(PostMortemPromptData{
Title: "Таймаут-задача",
Goal: "сделать",
Repos: []string{"calc"},
AC: "работает",
Status: storage.StatusTimeout,
Sessions: postMortemsText([]storage.Trace{tr}),
Sessions: postMortemsText([]storage.Trace{tr}, transcripts),
})
if err != nil {
t.Fatalf("render: %v", err)
}
for _, want := range []string{"Таймаут-задача", "timeout", "=== Агент: dev", "промпт dev", "вывод dev"} {
for _, want := range []string{
"Таймаут-задача", "timeout", "=== Агент: dev", "промпт dev", "вывод dev",
"Транскрипт сессии", "прочитал файл a.go", "правлю код",
} {
if !strings.Contains(prompt, want) {
t.Errorf("промпт не содержит %q", want)
}
}
}
// TestPostMortemsTextTranscriptUnavailable — сессия с session_id, для которой
// транскрипт не загружен, помечается как недоступный, а не падает.
func TestPostMortemsTextTranscriptUnavailable(t *testing.T) {
tr := storage.Trace{Agent: "reviewer", Status: storage.TraceFailed, SessionID: "sess-r"}
text := postMortemsText([]storage.Trace{tr}, nil)
if !strings.Contains(text, "транскрипт сессии недоступен") {
t.Errorf("нет пометки о недоступном транскрипте: %q", text)
}
}
// TestPostMortemsTextTranscriptEmpty — в сессии нет сообщений: агент не сделал
// ни одного шага — это пишется явно, чтобы постмортем не строил догадок.
func TestPostMortemsTextTranscriptEmpty(t *testing.T) {
tr := storage.Trace{Agent: "dev", Status: storage.TraceTimeout, SessionID: "sess-e"}
text := postMortemsText([]storage.Trace{tr}, map[string]string{"sess-e": ""})
if !strings.Contains(text, "не сделал ни одного шага") {
t.Errorf("нет пометки о пустой сессии: %q", text)
}
}
// TestFormatReviewVerdict — человекочитаемое описание вердикта ревьюера.
func TestFormatReviewVerdict(t *testing.T) {
tests := []struct {