test: починить тесты на Windows #6

Merged
Teraonious merged 1 commits from feat/3c8528d into main 2026-08-19 10:50:24 +05:00
6 changed files with 163 additions and 50 deletions
Showing only changes of commit b5fb583c90 - Show all commits

View File

@@ -25,6 +25,7 @@ func TestNew(t *testing.T) {
if err != nil { if err != nil {
t.Fatalf("New() err = %v", err) t.Fatalf("New() err = %v", err)
} }
defer a.Store.Close()
if a.Store == nil { if a.Store == nil {
t.Fatal("Store не создан") t.Fatal("Store не создан")
} }
@@ -87,6 +88,7 @@ func TestNew_RunCtxCancel(t *testing.T) {
if err != nil { if err != nil {
t.Fatalf("New() err = %v", err) t.Fatalf("New() err = %v", err)
} }
defer a.Store.Close()
ctx, cancel := context.WithCancel(context.Background()) ctx, cancel := context.WithCancel(context.Background())
cancel() // сразу отменяем cancel() // сразу отменяем
@@ -120,8 +122,8 @@ func TestNew_WorktreeCreated(t *testing.T) {
" token: \"test:token\"", " token: \"test:token\"",
" chat_id: \"12345\"", " chat_id: \"12345\"",
"paths:", "paths:",
" db: \"" + filepath.Join(tmp, "test.db") + "\"", " db: '" + filepath.Join(tmp, "test.db") + "'",
" worktree: \"" + wt + "\"", " worktree: '" + wt + "'",
"", "",
}, "\n") }, "\n")
if err := os.WriteFile(configPath, []byte(content), 0o600); err != nil { if err := os.WriteFile(configPath, []byte(content), 0o600); err != nil {
@@ -165,7 +167,7 @@ func TestNew_UpdateWiring(t *testing.T) {
" base_url: \"https://hub.example.com\"", " base_url: \"https://hub.example.com\"",
" token: \"cfg-update-token\"", " token: \"cfg-update-token\"",
"paths:", "paths:",
" db: \"" + dbPath + "\"", " db: '" + dbPath + "'",
"", // пустая строка в конце "", // пустая строка в конце
}, "\n") }, "\n")
if err := os.WriteFile(configPath, []byte(content), 0o600); err != nil { if err := os.WriteFile(configPath, []byte(content), 0o600); err != nil {
@@ -177,6 +179,7 @@ func TestNew_UpdateWiring(t *testing.T) {
if err != nil { if err != nil {
t.Fatalf("New() err = %v", err) t.Fatalf("New() err = %v", err)
} }
defer a.Store.Close()
if a.Updater == nil { if a.Updater == nil {
t.Fatal("Updater не создан") t.Fatal("Updater не создан")
} }
@@ -197,6 +200,7 @@ func TestNew_UpdateWiring(t *testing.T) {
if err != nil { if err != nil {
t.Fatalf("New() err = %v", err) t.Fatalf("New() err = %v", err)
} }
defer a2.Store.Close()
if a2.Updater.Token != "embedded-update-token" { if a2.Updater.Token != "embedded-update-token" {
t.Errorf("Token = %q, want embedded-update-token (вшитый приоритетнее)", a2.Updater.Token) t.Errorf("Token = %q, want embedded-update-token (вшитый приоритетнее)", a2.Updater.Token)
} }

View File

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

View File

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

View File

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

View File

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

View File

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