From b5fb583c904b0eb10ba9b8c3f6a11fca30e99065 Mon Sep 17 00:00:00 2001 From: "ki.sagidullin" Date: Wed, 19 Aug 2026 10:47:55 +0500 Subject: [PATCH] =?UTF-8?q?test:=20=D0=BF=D0=BE=D1=87=D0=B8=D0=BD=D0=B8?= =?UTF-8?q?=D1=82=D1=8C=20=D1=82=D0=B5=D1=81=D1=82=D1=8B=20=D0=BD=D0=B0=20?= =?UTF-8?q?Windows?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - E2E (app): e2eFakeAPI переведён на v2 HTTP API opencode (/api/*) с определением агента по тексту промпта; Router получает processed-счётчик и WaitProcessed, e2eChannel.deliver ждёт асинхронную обработку — убирает гонку «запрос сразу после deliver» и коллатеральный 'database is closed'. - app_test: одинарные YAML-кавычки для путей Windows (backslash-escape) + закрытие Store в TestNew/TestNew_RunCtxCancel/TestNew_UpdateWiring. - config_test: абсолютный путь строится с корнем тома (C:\...) и одинарными кавычками YAML. - opencode/server_test: fakeServeBin на Windows — .cmd с ping (#!/bin/sh не исполняется). - worker_test: TestWorkerSemaphore поллит до целевого статуса вместо фиксированных sleep (git на Windows медленнее). --- internal/app/app_test.go | 10 ++- internal/app/e2e_test.go | 116 +++++++++++++++++++++---------- internal/chat/router.go | 23 ++++++ internal/config/config_test.go | 11 +-- internal/opencode/server_test.go | 12 +++- internal/worker/worker_test.go | 41 +++++++++-- 6 files changed, 163 insertions(+), 50 deletions(-) diff --git a/internal/app/app_test.go b/internal/app/app_test.go index 900ef8b..e2f4d06 100644 --- a/internal/app/app_test.go +++ b/internal/app/app_test.go @@ -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) } diff --git a/internal/app/e2e_test.go b/internal/app/e2e_test.go index ab78626..973bd37 100644 --- a/internal/app/e2e_test.go +++ b/internal/app/e2e_test.go @@ -36,22 +36,35 @@ import ( ) // вердикты фейкового агента по имени. -var ( - e2eAgentVerdicts = map[string]string{ - "analyst": `{"phase":"propose","title":"Калькулятор","goal":"Сделать веб-калькулятор","repo":"calc","why":"Нужен для учёта","ac":"Работает + - * /","chat_reply":"Черновик готов."}`, - "dev": `done`, - "reviewer": `{"passed":true,"comments":[]}`, - } -) +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-) определяет -// агента и возвращает его вердикт как 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) } @@ -166,8 +200,9 @@ func e2eAssemble(t *testing.T) (*App, string, *e2eChannel) { // e2eChannel — минимальный fake-канал для перехвата исходящих // и доставки входящих через роутер (как реальный канал). type e2eChannel struct { - onMsg chat.Handler - sent []chat.Message + 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 { - c.onMsg(chat.Incoming{UserID: uid, Address: chat.Address("u://" + string(uid)), Msg: chat.Message{Text: text}, Channel: c}) + 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") } } diff --git a/internal/chat/router.go b/internal/chat/router.go index 95e355b..913cd0a 100644 --- a/internal/chat/router.go +++ b/internal/chat/router.go @@ -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 { diff --git a/internal/config/config_test.go b/internal/config/config_test.go index e77c7d4..d9bcba7 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -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 { diff --git a/internal/opencode/server_test.go b/internal/opencode/server_test.go index 086d00f..4c9ad5e 100644 --- a/internal/opencode/server_test.go +++ b/internal/opencode/server_test.go @@ -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" diff --git a/internal/worker/worker_test.go b/internal/worker/worker_test.go index 43db6e9..37f08dd 100644 --- a/internal/worker/worker_test.go +++ b/internal/worker/worker_test.go @@ -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 { -- 2.49.1