Merge pull request 'test: починить тесты на Windows' (#6) from feat/3c8528d into main
All checks were successful
CI / test (push) Successful in 35s
CI / build-and-package (amd64, linux) (push) Successful in 45s
CI / build-and-package (amd64, windows) (push) Successful in 42s

Reviewed-on: http://gitea.hal9000.home/kamelion/ratatoskr-go/pulls/6
This commit was merged in pull request #6.
This commit is contained in:
2026-08-19 10:50:23 +05:00
6 changed files with 163 additions and 50 deletions

View File

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

View File

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

View File

@@ -4,6 +4,8 @@ import (
"context"
"fmt"
"sync"
"sync/atomic"
"time"
)
// Router — единый диспетчер входящих из всех каналов и маршрутизатор исходящих.
@@ -24,6 +26,10 @@ type Router struct {
// long-poll цикл канала (Telegram) не блокируется на время долгого
// вызова аналитика и продолжает принимать новые сообщения.
incoming chan Incoming
// processed — число обработанных воркером событий (для синхронизации
// тестов с асинхронной очередью: WaitProcessed ждёт обработку события).
processed atomic.Int64
}
// NewRouter создаёт роутер. onUserMsg — колбэк обработки входящего.
@@ -46,9 +52,26 @@ func NewRouter(onUserMsg func(Incoming)) *Router {
func (r *Router) processLoop() {
for inc := range r.incoming {
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 регистрирует канал и подключает его к обработчику входящих.
// Возвращает ошибку только при пустом канале (nil).
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_CHAT_ID", "42")
absDB := filepath.Join(string(filepath.Separator), "data", "ratatoskr.db") // абсолютный для текущей ОС
absWt := filepath.Join(string(filepath.Separator), "worktrees")
// абсолютные пути «для текущей ОС»: на Windows слэш-относительный путь
// (\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:
username: "${TG_TOKEN}"
chat_id: "${TG_CHAT_ID}"
paths:
db: "` + absDB + `"
worktree: "` + absWt + `"
db: '` + absDB + `'
worktree: '` + absWt + `'
`
cfg, err := Load(writeCfg(t, yaml))
if err != nil {

View File

@@ -8,15 +8,25 @@ import (
"os"
"os/exec"
"path/filepath"
"runtime"
"strings"
"testing"
"time"
)
// fakeServeBin создаёт скрипт, имитирующий opencode serve: просто держит
// процесс живым (sleep), чтобы супервайзер мог им владеть и убивать его.
// процесс живым (sleep/ping), чтобы супервайзер мог им владеть и убивать его.
// На Windows используется .cmd (с #!/bin/sh нельзя — он не исполняется).
func fakeServeBin(t *testing.T, workdir string) string {
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")
script := `#!/bin/sh
echo "fake serve started"

View File

@@ -4,7 +4,6 @@ import (
"context"
"encoding/json"
"errors"
"fmt"
"os"
"os/exec"
"path/filepath"
@@ -809,13 +808,41 @@ func TestWorkerStartStop(t *testing.T) {
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) {
s := setupWorkerDB(t)
// создаём 2 ready-задачи
for i := 0; i < 2; i++ {
createReadyTask(t, s, fmt.Sprintf("task-%d", i))
}
task1 := createReadyTask(t, s, "task-0")
task2 := createReadyTask(t, s, "task-1")
w := &Worker{
Store: s,
@@ -829,12 +856,12 @@ func TestWorkerSemaphore(t *testing.T) {
w.sem = make(chan struct{}, 1)
w.sem <- struct{}{}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
// первый poll — запустит 1 задачу (макс. 1)
w.pollAndDispatch(ctx)
time.Sleep(200 * time.Millisecond)
waitTaskStatus(t, ctx, s, task1.ID, storage.StatusSuccess)
// 1 должна быть success, 1 — всё ещё approved
success, _ := s.ListTasks(ctx, storage.TaskFilter{Status: storage.StatusSuccess})
@@ -848,7 +875,7 @@ func TestWorkerSemaphore(t *testing.T) {
// первая завершилась и вернула токен в сем — можем диспатчить вторую
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})
if len(success) != 2 {