Compare commits
3 Commits
feat/ad025
...
0b31e18d85
| Author | SHA1 | Date | |
|---|---|---|---|
| 0b31e18d85 | |||
|
|
b5fb583c90 | ||
| 3c8528dbd9 |
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user