From 7eb3a0292c9c9188e34eb36a170c2cb2578fe330 Mon Sep 17 00:00:00 2001 From: Hermes Date: Tue, 18 Aug 2026 20:49:34 +0500 Subject: [PATCH] =?UTF-8?q?fix(opencode):=20=D1=83=D0=B1=D1=80=D0=B0=D1=82?= =?UTF-8?q?=D1=8C=20=D1=81=D0=BC=D0=B5=D1=88=D0=B5=D0=BD=D0=B8=D0=B5=20?= =?UTF-8?q?=D1=81=D0=BB=D0=BE=D1=91=D0=B2=20API=20=E2=80=94=20=D0=BF=D0=B5?= =?UTF-8?q?=D1=80=D0=B5=D0=B9=D1=82=D0=B8=20=D1=86=D0=B5=D0=BB=D0=B8=D0=BA?= =?UTF-8?q?=D0=BE=D0=BC=20=D0=BD=D0=B0=20experimental=20(/session)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Корень проблемы «не получаем результаты»: клиент смешивал два слоя opencode serve. CreateSession ходил на /api/session (v2, ждал {data.id}), Verdict — на /api/session/{id}/message?order=desc и ждал {data:[{type,content}]}, где поле content[].type/text физически отсутствует, поэтому вердикт никогда не находился и поллинг уходил в вечный таймаут. Abort и вовсе звал несуществующий /interrupt. Теперь весь код на experimental-слое, как сверено с sst/opencode (ветка dev): - CreateSession: POST /session → голая Session, id в .id. - Send: блокирующий POST /session/{id}/message, тело {parts:[{type:text,text}]}, вердикт из частей parts[].type=="text" ответа. Это и есть результат — метод Verdict и отдельный GET удалены. - textCount (прогресс): GET /session/{id}/message → голый массив [{info, parts}]. - Abort: POST /session/{id}/abort. Runner: блокирующий Send запускается в горутине (канал вердикта/ошибки), параллельно поллим textCount (рост text-частей сбрасывает idle-таймер). При idle/hard-таймауте или отмене контекста — Abort + cancel() Send-горутины → rc=-1. Send ходит через отдельный http.Client без жёсткого Timeout (управляется ctx), чтобы длинная генерация не обрывалась на 30s. Тесты/fakeAPIServer переведены на экспериментальный формат. Версия → 0.2.2. --- internal/app/app.go | 2 +- internal/app/e2e_test.go | 55 +++++----- internal/opencode/client.go | 174 +++++++++++++++++-------------- internal/opencode/client_test.go | 118 +++++++++++++-------- internal/opencode/runner.go | 117 ++++++++------------- internal/opencode/runner_test.go | 16 +-- 6 files changed, 251 insertions(+), 231 deletions(-) diff --git a/internal/app/app.go b/internal/app/app.go index fd436ae..93e77d1 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -49,7 +49,7 @@ const packageOwner = "kamelion" // не следует путать с build-идентификатором `main.version` (commit-), // который вшивается ldflag'ом и используется автообновлением. Здесь номер // поднимается вручную перед каждым релизом/публикацией новой сборки. -const Version = "0.2.1" +const Version = "0.2.2" // App — собранный конвейер. type App struct { diff --git a/internal/app/e2e_test.go b/internal/app/e2e_test.go index 9332144..ab78626 100644 --- a/internal/app/e2e_test.go +++ b/internal/app/e2e_test.go @@ -44,17 +44,24 @@ var ( } ) -// e2eFakeAPI поднимает фейковый opencode serve HTTP API v1.18 и возвращает URL. -// По title сессии (ratatoskr-) определяет агента и возвращает его вердикт -// как text-часть единственного assistant-сообщения. +// e2eFakeAPI поднимает фейковый opencode serve experimental HTTP API (пути +// БЕЗ /api) и возвращает URL. По title сессии (ratatoskr-) определяет +// агента и возвращает его вердикт как text-часть ответа на POST /message. func e2eFakeAPI(t *testing.T) string { t.Helper() var mu sync.Mutex sessions := map[string]string{} // id → agent + verdictFor := func(agent string) string { + if v, ok := e2eAgentVerdicts[agent]; ok { + return v + } + return "unknown agent" + } + h := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch { - case r.Method == http.MethodPost && r.URL.Path == "/api/session": + case r.Method == http.MethodPost && r.URL.Path == "/session": var req struct { Title string `json:"title"` } @@ -64,35 +71,27 @@ func e2eFakeAPI(t *testing.T) string { id := fmt.Sprintf("e2e-%d", len(sessions)+1) sessions[id] = agent mu.Unlock() - writeJSON(w, map[string]any{"data": map[string]any{"id": id}}) + // 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"): - w.WriteHeader(http.StatusOK) - _, _ = w.Write([]byte(`{"ok":true}`)) - - case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/wait"): - w.WriteHeader(http.StatusOK) - _, _ = w.Write([]byte(`{"ok":true}`)) - - case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/interrupt"): - w.WriteHeader(http.StatusOK) - _, _ = w.Write([]byte(`{"ok":true}`)) - - case r.Method == http.MethodGet && strings.HasSuffix(r.URL.Path, "/message"): - id := strings.TrimSuffix(strings.TrimPrefix(r.URL.Path, "/api/session/"), "/message") + // блокирующий ответ: вердикт как text-часть. + id := strings.TrimSuffix(strings.TrimPrefix(r.URL.Path, "/session/"), "/message") mu.Lock() agent := sessions[id] mu.Unlock() - verdict := "" - if v, ok := e2eAgentVerdicts[agent]; ok { - verdict = v - } else { - verdict = "unknown agent" - } - writeJSON(w, map[string]any{"data": []map[string]any{{ - "type": "assistant", - "content": []map[string]any{{"type": "text", "text": verdict}}, - }}}) + 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"): + // поллинг прогресса: голый массив [{info, parts}]. + 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.MethodPost && strings.HasSuffix(r.URL.Path, "/abort"): + writeJSON(w, map[string]any{"ok": true}) default: http.NotFound(w, r) diff --git a/internal/opencode/client.go b/internal/opencode/client.go index a35cf1b..6c7e795 100644 --- a/internal/opencode/client.go +++ b/internal/opencode/client.go @@ -13,26 +13,27 @@ import ( // Client — HTTP-взаимодействие с одним opencode serve (режим API). // -// Ходит по HTTP API opencode serve (v1.18+, префикс /api): -// - POST /api/session создать сессию → {data:{id}} -// - POST /session/:id/message отправить промпт {parts:[{type:"text"}]} -// - POST /api/session/:id/wait дождаться завершения ответа (блок) -// - GET /api/session/:id/message?order=desc → {data:[{...}]} история ответов -// - POST /api/session/:id/interrupt прервать выполняющийся ответ +// Ходит по experimental HTTP API opencode serve (пути БЕЗ префикса /api): +// - POST /session создать сессию → голая Session {id} +// - POST /session/{id}/message отправить промпт {parts:[{type:"text"}]} → +// блокирует и возвращает {info,parts}; вердикт из parts +// - GET /session/{id}/message история → голый массив [{info, parts}] (для прогресса) +// - POST /session/{id}/abort прервать выполняющийся ответ // -// Вердикт собирается из последнего assistant-сообщения: его content[] → текст -// тех частей, где type == "text". (В API-режиме текст в part.text — плоско, -// в отличие от NDJSON run, где он был вложен в part.part.text.) +// Вердикт собирается из parts[] ответа на POST /message: текст тех частей, +// где type == "text". type Client struct { BaseURL string // http://host:port (без завершающего слеша) Password string // basic auth (username "opencode") Debug bool // включать отладочные логи API-вызовов (log.level=debug) - http *http.Client + http *http.Client // для быстрых операций (create/messages/abort) + httpSend *http.Client // для блокирующего Send — без жёсткого таймаута, + // отменяется только через контекст (idle/hard) } // ClientErr — классы ошибок клиента. type ClientErr struct { - Op string // "connect" | "create" | "prompt" | "wait" | "messages" | "verdict" + Op string // "connect" | "create" | "prompt" | "messages" | "abort" Err error } @@ -43,11 +44,19 @@ func (c *Client) defaults() { if c.http == nil { c.http = &http.Client{Timeout: 30 * time.Second} } + if c.httpSend == nil { + c.httpSend = &http.Client{} + } } -// do выполняет запрос и возвращает тело ответа при 2xx. Иначе — ClientErr. -// op — метка операции (create/prompt/wait/messages/abort) для класса ошибки. +// do выполняет запрос через c.http (с таймаутом 30s) и возвращает тело при 2xx. func (c *Client) do(ctx context.Context, method, path, op string, body []byte) ([]byte, error) { + c.defaults() + return c.doHTTP(ctx, method, path, op, body, c.http) +} + +// doHTTP — общая реализация запроса; hc — клиент, которым выполняется запрос. +func (c *Client) doHTTP(ctx context.Context, method, path, op string, body []byte, hc *http.Client) ([]byte, error) { c.defaults() var rd io.Reader if body != nil { @@ -69,7 +78,7 @@ func (c *Client) do(ctx context.Context, method, path, op string, body []byte) ( log.Printf("opencode api %s request body: %s", op, truncateStr(string(body), 5000)) } } - resp, err := c.http.Do(req) + resp, err := hc.Do(req) if err != nil { return nil, &ClientErr{Op: "connect", Err: err} } @@ -97,96 +106,109 @@ func (c *Client) CreateSession(ctx context.Context, title string) (string, error body["title"] = title } b, _ := json.Marshal(body) - raw, err := c.do(ctx, http.MethodPost, "/api/session", "create", b) + raw, err := c.do(ctx, http.MethodPost, "/session", "create", b) if err != nil { return "", err } + // experimental: ответ — голая Session (без обёртки {data}). var out struct { - Data struct { - ID string `json:"id"` - } `json:"data"` + ID string `json:"id"` } if err := json.Unmarshal(raw, &out); err != nil { return "", &ClientErr{Op: "create", Err: fmt.Errorf("невалидный ответ: %v", err)} } - if out.Data.ID == "" { + if out.ID == "" { return "", &ClientErr{Op: "create", Err: fmt.Errorf("пустой id сессии")} } - return out.Data.ID, nil + return out.ID, nil } -// Prompt отправляет промпт в сессию и ждёт завершения ответа (блокирующий wait). -func (c *Client) Prompt(ctx context.Context, sessionID, prompt string) error { - if err := c.Send(ctx, sessionID, prompt); err != nil { - return err - } - // POST /wait блокирует, пока ответ сессии не завершится. - _, err := c.do(ctx, http.MethodPost, "/api/session/"+sessionID+"/wait", "wait", nil) - return err -} - -// Send отправляет промпт в сессию (не ждёт ответа — выполнение идёт в фоне). -// Используется Runner'ом для асинхронного ожидания через поллинг сообщений. -// Формат тела — v1.18: parts:[{type:"text"}], а не {prompt:{text}}. -func (c *Client) Send(ctx context.Context, sessionID, prompt string) error { +// Send отправляет промпт в сессию, БЛОКИРУЯСЬ до завершения ответа, и +// возвращает вердикт (текст text-частей из parts). Отмена — только через ctx +// (используется отдельный клиент без жёсткого таймаута; idle/hard в Runner'е +// отменяют контекст, что прерывает этот запрос). +func (c *Client) Send(ctx context.Context, sessionID, prompt string) (string, error) { payload := map[string]any{ "parts": []map[string]string{{"type": "text", "text": prompt}}, } b, _ := json.Marshal(payload) - _, err := c.do(ctx, http.MethodPost, "/session/"+sessionID+"/message", "prompt", b) - return err -} - -// Abort прерывает выполняющийся ответ сессии. -func (c *Client) Abort(ctx context.Context, sessionID string) error { - _, err := c.do(ctx, http.MethodPost, "/api/session/"+sessionID+"/interrupt", "abort", nil) - return err -} - -// sessionMessage — минимальная структура сообщения из GET /api/session/:id/message. -type sessionMessage struct { - Type string `json:"type"` // "assistant" | "user" | ... - Content []struct { - Type string `json:"type"` // "text" | "reasoning" | "tool" - Text string `json:"text"` - } `json:"content"` -} - -// Verdict возвращает текст последнего assistant-сообщения сессии (вердикт). -// Ошибка — класса ErrVerdict (ClientErr{Op:"verdict"}), если нет готового -// assistant-сообщения или в нём нет text-частей. -func (c *Client) Verdict(ctx context.Context, sessionID string) (string, error) { - // order=desc — самые свежие сообщения первыми; идём по ним в поисках - // первого незаконченного assistant-ответа. - raw, err := c.do(ctx, http.MethodGet, "/api/session/"+sessionID+"/message?order=desc&limit=50", "messages", nil) + c.defaults() + raw, err := c.doHTTP(ctx, http.MethodPost, "/session/"+sessionID+"/message", "prompt", b, c.httpSend) if err != nil { return "", err } var out struct { - Data []sessionMessage `json:"data"` + Parts []part `json:"parts"` } if err := json.Unmarshal(raw, &out); err != nil { - return "", &ClientErr{Op: "verdict", Err: fmt.Errorf("невалидный ответ: %v", err)} + return "", &ClientErr{Op: "prompt", Err: fmt.Errorf("невалидный ответ: %v", err)} } - var b bytes.Buffer - for _, m := range out.Data { - if m.Type != "assistant" { + var buf bytes.Buffer + for _, p := range out.Parts { + if p.Type == "text" && p.Text != "" { + if buf.Len() > 0 { + buf.WriteString("\n") + } + buf.WriteString(p.Text) + } + } + if buf.Len() == 0 { + return "", &ClientErr{Op: "prompt", Err: fmt.Errorf("нет text-части в ответе")} + } + return stripFence(buf.String()), nil +} + +// Abort прерывает выполняющийся ответ сессии. +func (c *Client) Abort(ctx context.Context, sessionID string) error { + _, err := c.do(ctx, http.MethodPost, "/session/"+sessionID+"/abort", "abort", nil) + return err +} + +// part — минимальная часть сообщения (из parts[]). +type part struct { + Type string `json:"type"` // "text" | "reasoning" | "tool" | ... + Text string `json:"text"` +} + +// message — элемент голого массива из GET /session/{id}/message. +type message struct { + Info struct { + Role string `json:"role"` // "assistant" | "user" | ... + } `json:"info"` + Parts []part `json:"parts"` +} + +// messages возвращает сырые сообщения сессии (для поллинга прогресса). +func (c *Client) messages(ctx context.Context, sessionID string) ([]message, error) { + raw, err := c.do(ctx, http.MethodGet, "/session/"+sessionID+"/message", "messages", nil) + if err != nil { + return nil, err + } + var out []message + if err := json.Unmarshal(raw, &out); err != nil { + return nil, err + } + return out, nil +} + +// textCount считает число text-частей в assistant-сообщениях (для progress). +func (c *Client) textCount(ctx context.Context, sessionID string) (int, error) { + msgs, err := c.messages(ctx, sessionID) + if err != nil { + return 0, err + } + n := 0 + for _, m := range msgs { + if m.Info.Role != "assistant" { continue } - b.Reset() - for _, p := range m.Content { + for _, p := range m.Parts { if p.Type == "text" && p.Text != "" { - if b.Len() > 0 { - b.WriteString("\n") - } - b.WriteString(p.Text) + n++ } } - if b.Len() > 0 { - return stripFence(b.String()), nil - } } - return "", &ClientErr{Op: "verdict", Err: fmt.Errorf("нет assistant-сообщения с text-частью")} + return n, nil } func truncateStr(s string, n int) string { @@ -194,4 +216,4 @@ func truncateStr(s string, n int) string { return s } return s[:n] + "..." -} \ No newline at end of file +} diff --git a/internal/opencode/client_test.go b/internal/opencode/client_test.go index 05f0b1a..e231f87 100644 --- a/internal/opencode/client_test.go +++ b/internal/opencode/client_test.go @@ -7,41 +7,69 @@ import ( "net/http" "net/http/httptest" "testing" + "time" ) -// fakeAPIServer — минимальный фейк opencode serve HTTP API v1.18. +// fakeAPIServer — минимальный фейк opencode serve experimental HTTP API +// (пути БЕЗ префикса /api). type fakeAPIServer struct { - messages []sessionMessage - failCreate bool - failVerify bool + messages []message + failCreate bool + verdictParts []part // ответ на POST /session/{id}/message (вердикт) + blockPrompt bool // POST /message блокируется до отмены ctx (эмуляция зависания) } func (f *fakeAPIServer) handler() http.Handler { mux := http.NewServeMux() - mux.HandleFunc("/api/session", func(w http.ResponseWriter, r *http.Request) { + mux.HandleFunc("/session", func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + w.WriteHeader(http.StatusMethodNotAllowed) + return + } if f.failCreate { http.Error(w, "boom", http.StatusInternalServerError) return } - writeJSON(w, map[string]any{"data": map[string]any{"id": "sess-fake"}}) + // experimental: голая Session (без обёртки {data}). + writeJSON(w, map[string]any{"id": "sess-fake"}) }) - mux.HandleFunc("/api/session/{id}/wait", func(w http.ResponseWriter, _ *http.Request) { - writeJSON(w, map[string]any{}) - }) - mux.HandleFunc("/api/session/{id}/interrupt", func(w http.ResponseWriter, _ *http.Request) { - writeJSON(w, map[string]any{}) - }) - // POST отправка промпта (v1.18) — путь БЕЗ /api. - mux.HandleFunc("/session/{id}/message", func(w http.ResponseWriter, r *http.Request) { - writeJSON(w, map[string]any{"data": map[string]any{}}) - }) - // GET чтение сообщений — путь с /api. - mux.HandleFunc("/api/session/{id}/message", func(w http.ResponseWriter, r *http.Request) { - if len(f.messages) == 0 { - writeJSON(w, map[string]any{"data": []sessionMessage{}}) + mux.HandleFunc("/session/{id}/abort", func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + w.WriteHeader(http.StatusMethodNotAllowed) return } - writeJSON(w, map[string]any{"data": f.messages}) + writeJSON(w, map[string]any{}) + }) + mux.HandleFunc("/session/{id}/message", func(w http.ResponseWriter, r *http.Request) { + switch r.Method { + case http.MethodPost: + if f.blockPrompt { + // Эмуляция «зависшего» агента: ответ приходит позже idle-таймаута, + // но handler всё равно завершится, чтобы не блокировать shutdown. + select { + case <-r.Context().Done(): + case <-time.After(2 * time.Second): + } + w.WriteHeader(http.StatusRequestTimeout) + return + } + // блокирующий ответ: {info, parts}, где вердикт — text-части. + info := map[string]any{"role": "assistant"} + parts := f.verdictParts + if parts == nil { + parts = []part{} + } + writeJSON(w, map[string]any{"info": info, "parts": parts}) + case http.MethodGet: + // голый массив [{info, parts}]. + if f.messages == nil { + writeJSON(w, []message{}) + return + } + writeJSON(w, f.messages) + default: + w.WriteHeader(http.StatusMethodNotAllowed) + } }) return mux } @@ -51,17 +79,7 @@ func writeJSON(w http.ResponseWriter, v any) { _ = json.NewEncoder(w).Encode(v) } -// textPart — анонимная text-часть assistant-сообщения. -func textPart(s string) struct { - Type string `json:"type"` - Text string `json:"text"` -} { - return struct { - Type string `json:"type"` - Text string `json:"text"` - }{Type: "text", Text: s} -} - +// fakeClient — клиент к фейк-серверу. func fakeClient(t *testing.T, f *fakeAPIServer) *Client { t.Helper() ts := httptest.NewServer(f.handler()) @@ -87,30 +105,44 @@ func TestClient_CreateSessionFail(t *testing.T) { } } -func TestClient_Verdict(t *testing.T) { +func TestClient_Send(t *testing.T) { c := fakeClient(t, &fakeAPIServer{ - messages: []sessionMessage{{Type: "assistant", Content: []struct { - Type string `json:"type"` - Text string `json:"text"` - }{textPart(`{"phase":"ready"}`)}}}, + verdictParts: []part{{Type: "text", Text: `{"phase":"ready"}`}}, }) - vd, err := c.Verdict(context.Background(), "sess-fake") + vd, err := c.Send(context.Background(), "sess-fake", "почини x") if err != nil { - t.Fatalf("Verdict err: %v", err) + t.Fatalf("Send err: %v", err) } if vd != `{"phase":"ready"}` { t.Errorf("verdict = %q, want вердикт модели", vd) } } -func TestClient_VerdictNoText(t *testing.T) { - c := fakeClient(t, &fakeAPIServer{}) // нет assistant-сообщения с text - if _, err := c.Verdict(context.Background(), "sess-fake"); err == nil { - t.Fatal("Verdict должен упасть, когда нет text-части") +func TestClient_SendNoText(t *testing.T) { + c := fakeClient(t, &fakeAPIServer{}) // нет text-части в ответе + if _, err := c.Send(context.Background(), "sess-fake", "почини x"); err == nil { + t.Fatal("Send должен упасть, когда нет text-части") } else { var ce *ClientErr if !errors.As(err, &ce) { t.Errorf("ожидался *ClientErr, got %T", err) } } +} + +func TestClient_textCount(t *testing.T) { + f := &fakeAPIServer{messages: []message{{ + Info: struct { + Role string `json:"role"` + }{Role: "assistant"}, + Parts: []part{{Type: "text", Text: "a"}, {Type: "reasoning", Text: "x"}}, + }}} + c := fakeClient(t, f) + n, err := c.textCount(context.Background(), "sess-fake") + if err != nil { + t.Fatalf("textCount err: %v", err) + } + if n != 1 { + t.Errorf("textCount = %d, want 1 (одна text-часть в assistant)", n) + } } \ No newline at end of file diff --git a/internal/opencode/runner.go b/internal/opencode/runner.go index efb5854..37d9d20 100644 --- a/internal/opencode/runner.go +++ b/internal/opencode/runner.go @@ -2,7 +2,6 @@ package opencode import ( "context" - "encoding/json" "errors" "fmt" "io" @@ -84,17 +83,27 @@ func (r *Runner) Run(ctx context.Context, prompt, cwd, agent, sessionID string) r.logf("opencode(%s) session=%s на %s", agent, sid, srv.Addr()) } - // Отправляем промпт (неблокирующий — сервер начинает выполнение). - if err := c.Send(ctx, sid, prompt); err != nil { - return nil, fmt.Errorf("opencode: prompt: %w", err) - } - - return r.awaitVerdict(ctx, c, sid, agent) + // Отправляем промпт (блокирующий Send в горутине; вердикт придёт из него), + // параллельно поллим прогресс и контролируем idle/hard таймауты. + return r.awaitVerdict(ctx, c, sid, agent, prompt) } -// awaitVerdict поллит сообщения сессии, пока не появится готовый text-вердикт -// от assistant, либо не истечёт idle/hard таймаут (тогда Abort + rc=-1). -func (r *Runner) awaitVerdict(ctx context.Context, c *Client, sid, agent string) (*Result, error) { +// awaitVerdict запускает блокирующий Send и параллельно поллит прогресс +// (рост числа text-частей = агент жив, сбрасывает idle). Возвращается вердикт +// из ответа Send, либо rc=-1 при idle/hard таймауте (тогда Abort + отмена ctx). +func (r *Runner) awaitVerdict(ctx context.Context, c *Client, sid, agent, prompt string) (*Result, error) { + sendCtx, cancel := context.WithCancel(ctx) + defer cancel() + type sendOut struct { + vd string + err error + } + sendCh := make(chan sendOut, 1) + go func() { + vd, err := c.Send(sendCtx, sid, prompt) + sendCh <- sendOut{vd: vd, err: err} + }() + // Прогресс = сумма text-частей во всех assistant-сообщениях сессии. Рост // сбрасывает idle-таймер (LLM стримит = жив). var mu sync.Mutex @@ -102,13 +111,18 @@ func (r *Runner) awaitVerdict(ctx context.Context, c *Client, sid, agent string) lastProgress := time.Now() launch := time.Now() + abortAnd := func(rc int, why string) (*Result, error) { + if err := c.Abort(ctx, sid); err != nil { + r.logf("opencode(%s) abort %s: %v", agent, why, err) + } + cancel() + return &Result{RC: rc, Stdout: "", SessionID: sid}, nil + } + for { if ctx.Err() != nil { - if err := c.Abort(ctx, sid); err != nil { - r.logf("opencode(%s) abort (ctx): %v", agent, err) - } r.logf("opencode(%s) ctx cancelled — обрыв (rc=-1)", agent) - return &Result{RC: -1, Stdout: "", SessionID: sid}, nil + return abortAnd(-1, "ctx") } count, _ := c.textCount(ctx, sid) @@ -119,80 +133,39 @@ func (r *Runner) awaitVerdict(ctx context.Context, c *Client, sid, agent string) } mu.Unlock() - // Пробуем вердикт (дешёвый GET). Если готов — выходим. - vd, vErr := c.Verdict(ctx, sid) - if vErr == nil && vd != "" { - r.logf("opencode(%s) вердикт готов (%d байт)", agent, len(vd)) - return &Result{RC: 0, Stdout: vd, SessionID: sid}, nil - } - // Сервер недоступен — фатально (не таймаут). Но если это следствие - // отмены контекста (cancel прилетел прямо во время запроса) — это обрыв - // по контексту, а не недоступность сервера: вернём RC=-1 на следующей - // итерации (проверка ctx.Err() вверху) и не подменяем класс ошибки. - var ce *ClientErr - if errors.As(vErr, &ce) && ce.Op == "connect" && ctx.Err() == nil { - return nil, fmt.Errorf("opencode: %w", vErr) - } - now := time.Now() if now.Sub(lastProgress) > r.IdleTimeout { r.logf("opencode(%s) idle %.0fs — abort", agent, r.IdleTimeout.Seconds()) - if err := c.Abort(ctx, sid); err != nil { - r.logf("opencode(%s) abort idle: %v", agent, err) - } - return &Result{RC: -1, Stdout: "", SessionID: sid}, nil + return abortAnd(-1, "idle") } // hard — общий бюджет от старта запуска. if now.Sub(launch) > r.HardTimeout { r.logf("opencode(%s) hard timeout %.0fs — abort", agent, r.HardTimeout.Seconds()) - if err := c.Abort(ctx, sid); err != nil { - r.logf("opencode(%s) abort hard: %v", agent, err) - } - return &Result{RC: -1, Stdout: "", SessionID: sid}, nil + return abortAnd(-1, "hard") } select { + case out := <-sendCh: + // Send завершился. Ошибка — connect (сервер недоступен) и ctx жив → + // фатально, не таймаут. Если ctx уже отменён — это обрыв, а не ошибка. + if out.err != nil { + var ce *ClientErr + if errors.As(out.err, &ce) && ce.Op == "connect" && ctx.Err() == nil { + return nil, fmt.Errorf("opencode: %w", out.err) + } + if ctx.Err() != nil { + return abortAnd(-1, "ctx") + } + return nil, out.err + } + r.logf("opencode(%s) вердикт готов (%d байт)", agent, len(out.vd)) + return &Result{RC: 0, Stdout: out.vd, SessionID: sid}, nil case <-time.After(r.PollInterval): case <-ctx.Done(): } } } -// textCount считает число text-частей в assistant-сообщениях (для progress). -func (c *Client) textCount(ctx context.Context, sessionID string) (int, error) { - msgs, err := c.messages(ctx, sessionID) - if err != nil { - return 0, err - } - n := 0 - for _, m := range msgs { - if m.Type != "assistant" { - continue - } - for _, p := range m.Content { - if p.Type == "text" && p.Text != "" { - n++ - } - } - } - return n, nil -} - -// messages возвращает сырые сообщения сессии (для поллинга прогресса). -func (c *Client) messages(ctx context.Context, sessionID string) ([]sessionMessage, error) { - raw, err := c.do(ctx, "GET", "/api/session/"+sessionID+"/message?order=asc&limit=200", "messages", nil) - if err != nil { - return nil, err - } - var out struct { - Data []sessionMessage `json:"data"` - } - if err := json.Unmarshal(raw, &out); err != nil { - return nil, err - } - return out.Data, nil -} - // ResumeDev — запуск dev-агента с resume-fallback. Если resume (sessionID) // падает с rc!=0 — повторяем ОДИН раз свежей сессией в том же каталоге. // rc=-1 (обрыв по таймауту) НЕ триггерит fallback. diff --git a/internal/opencode/runner_test.go b/internal/opencode/runner_test.go index 535b67b..d890f62 100644 --- a/internal/opencode/runner_test.go +++ b/internal/opencode/runner_test.go @@ -30,10 +30,7 @@ func httptestURL(t *testing.T, f *fakeAPIServer) string { func TestRun_Success(t *testing.T) { dir := t.TempDir() f := &fakeAPIServer{ - messages: []sessionMessage{{Type: "assistant", Content: []struct { - Type string `json:"type"` - Text string `json:"text"` - }{textPart("done")}}}, + verdictParts: []part{{Type: "text", Text: "done"}}, } p, _ := fakePool(t, f, dir) @@ -55,8 +52,8 @@ func TestRun_Success(t *testing.T) { func TestRun_IdleTimeout(t *testing.T) { dir := t.TempDir() - // сервер никогда не отдаёт text → всегда неготов, прогресс не растёт - f := &fakeAPIServer{} + // prompt блокируется (агент «завис»), прогресс не растёт → idle abort + f := &fakeAPIServer{blockPrompt: true} p, _ := fakePool(t, f, dir) r := &Runner{Pool: p, IdleTimeout: 30 * time.Millisecond, @@ -72,7 +69,7 @@ func TestRun_IdleTimeout(t *testing.T) { func TestRun_ContextCancel(t *testing.T) { dir := t.TempDir() - f := &fakeAPIServer{} + f := &fakeAPIServer{blockPrompt: true} p, _ := fakePool(t, f, dir) ctx, cancel := context.WithCancel(context.Background()) @@ -102,10 +99,7 @@ func TestRun_ContextCancel(t *testing.T) { func TestResumeDev_NoFallbackOnSuccess(t *testing.T) { dir := t.TempDir() f := &fakeAPIServer{ - messages: []sessionMessage{{Type: "assistant", Content: []struct { - Type string `json:"type"` - Text string `json:"text"` - }{textPart("ok")}}}, + verdictParts: []part{{Type: "text", Text: "ok"}}, } p, _ := fakePool(t, f, dir)