fix(opencode): убрать смешение слоёв API — перейти целиком на experimental (/session)
Корень проблемы «не получаем результаты»: клиент смешивал два слоя 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.
This commit is contained in:
@@ -49,7 +49,7 @@ const packageOwner = "kamelion"
|
|||||||
// не следует путать с build-идентификатором `main.version` (commit-<sha7>),
|
// не следует путать с build-идентификатором `main.version` (commit-<sha7>),
|
||||||
// который вшивается ldflag'ом и используется автообновлением. Здесь номер
|
// который вшивается ldflag'ом и используется автообновлением. Здесь номер
|
||||||
// поднимается вручную перед каждым релизом/публикацией новой сборки.
|
// поднимается вручную перед каждым релизом/публикацией новой сборки.
|
||||||
const Version = "0.2.1"
|
const Version = "0.2.2"
|
||||||
|
|
||||||
// App — собранный конвейер.
|
// App — собранный конвейер.
|
||||||
type App struct {
|
type App struct {
|
||||||
|
|||||||
@@ -44,17 +44,24 @@ var (
|
|||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
// e2eFakeAPI поднимает фейковый opencode serve HTTP API v1.18 и возвращает URL.
|
// e2eFakeAPI поднимает фейковый opencode serve experimental HTTP API (пути
|
||||||
// По title сессии (ratatoskr-<agent>) определяет агента и возвращает его вердикт
|
// БЕЗ /api) и возвращает URL. По title сессии (ratatoskr-<agent>) определяет
|
||||||
// как text-часть единственного assistant-сообщения.
|
// агента и возвращает его вердикт как text-часть ответа на POST /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
|
||||||
|
|
||||||
|
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) {
|
h := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
switch {
|
switch {
|
||||||
case r.Method == http.MethodPost && r.URL.Path == "/api/session":
|
case r.Method == http.MethodPost && r.URL.Path == "/session":
|
||||||
var req struct {
|
var req struct {
|
||||||
Title string `json:"title"`
|
Title string `json:"title"`
|
||||||
}
|
}
|
||||||
@@ -64,35 +71,27 @@ func e2eFakeAPI(t *testing.T) string {
|
|||||||
id := fmt.Sprintf("e2e-%d", len(sessions)+1)
|
id := fmt.Sprintf("e2e-%d", len(sessions)+1)
|
||||||
sessions[id] = agent
|
sessions[id] = agent
|
||||||
mu.Unlock()
|
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"):
|
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/message"):
|
||||||
w.WriteHeader(http.StatusOK)
|
// блокирующий ответ: вердикт как text-часть.
|
||||||
_, _ = w.Write([]byte(`{"ok":true}`))
|
id := strings.TrimSuffix(strings.TrimPrefix(r.URL.Path, "/session/"), "/message")
|
||||||
|
|
||||||
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")
|
|
||||||
mu.Lock()
|
mu.Lock()
|
||||||
agent := sessions[id]
|
agent := sessions[id]
|
||||||
mu.Unlock()
|
mu.Unlock()
|
||||||
verdict := ""
|
writeJSON(w, map[string]any{"info": map[string]any{"role": "assistant"}, "parts": []map[string]any{{"type": "text", "text": verdictFor(agent)}}})
|
||||||
if v, ok := e2eAgentVerdicts[agent]; ok {
|
|
||||||
verdict = v
|
case r.Method == http.MethodGet && strings.HasSuffix(r.URL.Path, "/message"):
|
||||||
} else {
|
// поллинг прогресса: голый массив [{info, parts}].
|
||||||
verdict = "unknown agent"
|
id := strings.TrimSuffix(strings.TrimPrefix(r.URL.Path, "/session/"), "/message")
|
||||||
}
|
mu.Lock()
|
||||||
writeJSON(w, map[string]any{"data": []map[string]any{{
|
agent := sessions[id]
|
||||||
"type": "assistant",
|
mu.Unlock()
|
||||||
"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.MethodPost && strings.HasSuffix(r.URL.Path, "/abort"):
|
||||||
|
writeJSON(w, map[string]any{"ok": true})
|
||||||
|
|
||||||
default:
|
default:
|
||||||
http.NotFound(w, r)
|
http.NotFound(w, r)
|
||||||
|
|||||||
@@ -13,26 +13,27 @@ import (
|
|||||||
|
|
||||||
// Client — HTTP-взаимодействие с одним opencode serve (режим API).
|
// Client — HTTP-взаимодействие с одним opencode serve (режим API).
|
||||||
//
|
//
|
||||||
// Ходит по HTTP API opencode serve (v1.18+, префикс /api):
|
// Ходит по experimental HTTP API opencode serve (пути БЕЗ префикса /api):
|
||||||
// - POST /api/session создать сессию → {data:{id}}
|
// - POST /session создать сессию → голая Session {id}
|
||||||
// - POST /session/:id/message отправить промпт {parts:[{type:"text"}]}
|
// - POST /session/{id}/message отправить промпт {parts:[{type:"text"}]} →
|
||||||
// - POST /api/session/:id/wait дождаться завершения ответа (блок)
|
// блокирует и возвращает {info,parts}; вердикт из parts
|
||||||
// - GET /api/session/:id/message?order=desc → {data:[{...}]} история ответов
|
// - GET /session/{id}/message история → голый массив [{info, parts}] (для прогресса)
|
||||||
// - POST /api/session/:id/interrupt прервать выполняющийся ответ
|
// - POST /session/{id}/abort прервать выполняющийся ответ
|
||||||
//
|
//
|
||||||
// Вердикт собирается из последнего assistant-сообщения: его content[] → текст
|
// Вердикт собирается из parts[] ответа на POST /message: текст тех частей,
|
||||||
// тех частей, где type == "text". (В API-режиме текст в part.text — плоско,
|
// где type == "text".
|
||||||
// в отличие от NDJSON run, где он был вложен в part.part.text.)
|
|
||||||
type Client struct {
|
type Client struct {
|
||||||
BaseURL string // http://host:port (без завершающего слеша)
|
BaseURL string // http://host:port (без завершающего слеша)
|
||||||
Password string // basic auth (username "opencode")
|
Password string // basic auth (username "opencode")
|
||||||
Debug bool // включать отладочные логи API-вызовов (log.level=debug)
|
Debug bool // включать отладочные логи API-вызовов (log.level=debug)
|
||||||
http *http.Client
|
http *http.Client // для быстрых операций (create/messages/abort)
|
||||||
|
httpSend *http.Client // для блокирующего Send — без жёсткого таймаута,
|
||||||
|
// отменяется только через контекст (idle/hard)
|
||||||
}
|
}
|
||||||
|
|
||||||
// ClientErr — классы ошибок клиента.
|
// ClientErr — классы ошибок клиента.
|
||||||
type ClientErr struct {
|
type ClientErr struct {
|
||||||
Op string // "connect" | "create" | "prompt" | "wait" | "messages" | "verdict"
|
Op string // "connect" | "create" | "prompt" | "messages" | "abort"
|
||||||
Err error
|
Err error
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -43,11 +44,19 @@ func (c *Client) defaults() {
|
|||||||
if c.http == nil {
|
if c.http == nil {
|
||||||
c.http = &http.Client{Timeout: 30 * time.Second}
|
c.http = &http.Client{Timeout: 30 * time.Second}
|
||||||
}
|
}
|
||||||
|
if c.httpSend == nil {
|
||||||
|
c.httpSend = &http.Client{}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// do выполняет запрос и возвращает тело ответа при 2xx. Иначе — ClientErr.
|
// do выполняет запрос через c.http (с таймаутом 30s) и возвращает тело при 2xx.
|
||||||
// op — метка операции (create/prompt/wait/messages/abort) для класса ошибки.
|
|
||||||
func (c *Client) do(ctx context.Context, method, path, op string, body []byte) ([]byte, error) {
|
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()
|
c.defaults()
|
||||||
var rd io.Reader
|
var rd io.Reader
|
||||||
if body != nil {
|
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))
|
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 {
|
if err != nil {
|
||||||
return nil, &ClientErr{Op: "connect", Err: err}
|
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
|
body["title"] = title
|
||||||
}
|
}
|
||||||
b, _ := json.Marshal(body)
|
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 {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
|
// experimental: ответ — голая Session (без обёртки {data}).
|
||||||
var out struct {
|
var out struct {
|
||||||
Data struct {
|
|
||||||
ID string `json:"id"`
|
ID string `json:"id"`
|
||||||
} `json:"data"`
|
|
||||||
}
|
}
|
||||||
if err := json.Unmarshal(raw, &out); err != nil {
|
if err := json.Unmarshal(raw, &out); err != nil {
|
||||||
return "", &ClientErr{Op: "create", Err: fmt.Errorf("невалидный ответ: %v", err)}
|
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 "", &ClientErr{Op: "create", Err: fmt.Errorf("пустой id сессии")}
|
||||||
}
|
}
|
||||||
return out.Data.ID, nil
|
return out.ID, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Prompt отправляет промпт в сессию и ждёт завершения ответа (блокирующий wait).
|
// Send отправляет промпт в сессию, БЛОКИРУЯСЬ до завершения ответа, и
|
||||||
func (c *Client) Prompt(ctx context.Context, sessionID, prompt string) error {
|
// возвращает вердикт (текст text-частей из parts). Отмена — только через ctx
|
||||||
if err := c.Send(ctx, sessionID, prompt); err != nil {
|
// (используется отдельный клиент без жёсткого таймаута; idle/hard в Runner'е
|
||||||
return err
|
// отменяют контекст, что прерывает этот запрос).
|
||||||
}
|
func (c *Client) Send(ctx context.Context, sessionID, prompt string) (string, error) {
|
||||||
// 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 {
|
|
||||||
payload := map[string]any{
|
payload := map[string]any{
|
||||||
"parts": []map[string]string{{"type": "text", "text": prompt}},
|
"parts": []map[string]string{{"type": "text", "text": prompt}},
|
||||||
}
|
}
|
||||||
b, _ := json.Marshal(payload)
|
b, _ := json.Marshal(payload)
|
||||||
_, err := c.do(ctx, http.MethodPost, "/session/"+sessionID+"/message", "prompt", b)
|
c.defaults()
|
||||||
return err
|
raw, err := c.doHTTP(ctx, http.MethodPost, "/session/"+sessionID+"/message", "prompt", b, c.httpSend)
|
||||||
}
|
|
||||||
|
|
||||||
// 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)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
var out struct {
|
var out struct {
|
||||||
Data []sessionMessage `json:"data"`
|
Parts []part `json:"parts"`
|
||||||
}
|
}
|
||||||
if err := json.Unmarshal(raw, &out); err != nil {
|
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
|
var buf bytes.Buffer
|
||||||
for _, m := range out.Data {
|
for _, p := range out.Parts {
|
||||||
if m.Type != "assistant" {
|
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
|
continue
|
||||||
}
|
}
|
||||||
b.Reset()
|
for _, p := range m.Parts {
|
||||||
for _, p := range m.Content {
|
|
||||||
if p.Type == "text" && p.Text != "" {
|
if p.Type == "text" && p.Text != "" {
|
||||||
if b.Len() > 0 {
|
n++
|
||||||
b.WriteString("\n")
|
|
||||||
}
|
|
||||||
b.WriteString(p.Text)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if b.Len() > 0 {
|
|
||||||
return stripFence(b.String()), nil
|
|
||||||
}
|
}
|
||||||
}
|
return n, nil
|
||||||
return "", &ClientErr{Op: "verdict", Err: fmt.Errorf("нет assistant-сообщения с text-частью")}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func truncateStr(s string, n int) string {
|
func truncateStr(s string, n int) string {
|
||||||
|
|||||||
@@ -7,41 +7,69 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
// fakeAPIServer — минимальный фейк opencode serve HTTP API v1.18.
|
// fakeAPIServer — минимальный фейк opencode serve experimental HTTP API
|
||||||
|
// (пути БЕЗ префикса /api).
|
||||||
type fakeAPIServer struct {
|
type fakeAPIServer struct {
|
||||||
messages []sessionMessage
|
messages []message
|
||||||
failCreate bool
|
failCreate bool
|
||||||
failVerify bool
|
verdictParts []part // ответ на POST /session/{id}/message (вердикт)
|
||||||
|
blockPrompt bool // POST /message блокируется до отмены ctx (эмуляция зависания)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *fakeAPIServer) handler() http.Handler {
|
func (f *fakeAPIServer) handler() http.Handler {
|
||||||
mux := http.NewServeMux()
|
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 {
|
if f.failCreate {
|
||||||
http.Error(w, "boom", http.StatusInternalServerError)
|
http.Error(w, "boom", http.StatusInternalServerError)
|
||||||
return
|
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) {
|
mux.HandleFunc("/session/{id}/abort", func(w http.ResponseWriter, r *http.Request) {
|
||||||
writeJSON(w, map[string]any{})
|
if r.Method != http.MethodPost {
|
||||||
})
|
w.WriteHeader(http.StatusMethodNotAllowed)
|
||||||
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{}})
|
|
||||||
return
|
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
|
return mux
|
||||||
}
|
}
|
||||||
@@ -51,17 +79,7 @@ func writeJSON(w http.ResponseWriter, v any) {
|
|||||||
_ = json.NewEncoder(w).Encode(v)
|
_ = json.NewEncoder(w).Encode(v)
|
||||||
}
|
}
|
||||||
|
|
||||||
// textPart — анонимная text-часть assistant-сообщения.
|
// fakeClient — клиент к фейк-серверу.
|
||||||
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}
|
|
||||||
}
|
|
||||||
|
|
||||||
func fakeClient(t *testing.T, f *fakeAPIServer) *Client {
|
func fakeClient(t *testing.T, f *fakeAPIServer) *Client {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
ts := httptest.NewServer(f.handler())
|
ts := httptest.NewServer(f.handler())
|
||||||
@@ -87,26 +105,23 @@ func TestClient_CreateSessionFail(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestClient_Verdict(t *testing.T) {
|
func TestClient_Send(t *testing.T) {
|
||||||
c := fakeClient(t, &fakeAPIServer{
|
c := fakeClient(t, &fakeAPIServer{
|
||||||
messages: []sessionMessage{{Type: "assistant", Content: []struct {
|
verdictParts: []part{{Type: "text", Text: `{"phase":"ready"}`}},
|
||||||
Type string `json:"type"`
|
|
||||||
Text string `json:"text"`
|
|
||||||
}{textPart(`{"phase":"ready"}`)}}},
|
|
||||||
})
|
})
|
||||||
vd, err := c.Verdict(context.Background(), "sess-fake")
|
vd, err := c.Send(context.Background(), "sess-fake", "почини x")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("Verdict err: %v", err)
|
t.Fatalf("Send err: %v", err)
|
||||||
}
|
}
|
||||||
if vd != `{"phase":"ready"}` {
|
if vd != `{"phase":"ready"}` {
|
||||||
t.Errorf("verdict = %q, want вердикт модели", vd)
|
t.Errorf("verdict = %q, want вердикт модели", vd)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestClient_VerdictNoText(t *testing.T) {
|
func TestClient_SendNoText(t *testing.T) {
|
||||||
c := fakeClient(t, &fakeAPIServer{}) // нет assistant-сообщения с text
|
c := fakeClient(t, &fakeAPIServer{}) // нет text-части в ответе
|
||||||
if _, err := c.Verdict(context.Background(), "sess-fake"); err == nil {
|
if _, err := c.Send(context.Background(), "sess-fake", "почини x"); err == nil {
|
||||||
t.Fatal("Verdict должен упасть, когда нет text-части")
|
t.Fatal("Send должен упасть, когда нет text-части")
|
||||||
} else {
|
} else {
|
||||||
var ce *ClientErr
|
var ce *ClientErr
|
||||||
if !errors.As(err, &ce) {
|
if !errors.As(err, &ce) {
|
||||||
@@ -114,3 +129,20 @@ func TestClient_VerdictNoText(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -2,7 +2,6 @@ package opencode
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"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())
|
r.logf("opencode(%s) session=%s на %s", agent, sid, srv.Addr())
|
||||||
}
|
}
|
||||||
|
|
||||||
// Отправляем промпт (неблокирующий — сервер начинает выполнение).
|
// Отправляем промпт (блокирующий Send в горутине; вердикт придёт из него),
|
||||||
if err := c.Send(ctx, sid, prompt); err != nil {
|
// параллельно поллим прогресс и контролируем idle/hard таймауты.
|
||||||
return nil, fmt.Errorf("opencode: prompt: %w", err)
|
return r.awaitVerdict(ctx, c, sid, agent, prompt)
|
||||||
}
|
|
||||||
|
|
||||||
return r.awaitVerdict(ctx, c, sid, agent)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// awaitVerdict поллит сообщения сессии, пока не появится готовый text-вердикт
|
// awaitVerdict запускает блокирующий Send и параллельно поллит прогресс
|
||||||
// от assistant, либо не истечёт idle/hard таймаут (тогда Abort + rc=-1).
|
// (рост числа text-частей = агент жив, сбрасывает idle). Возвращается вердикт
|
||||||
func (r *Runner) awaitVerdict(ctx context.Context, c *Client, sid, agent string) (*Result, error) {
|
// из ответа 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-сообщениях сессии. Рост
|
// Прогресс = сумма text-частей во всех assistant-сообщениях сессии. Рост
|
||||||
// сбрасывает idle-таймер (LLM стримит = жив).
|
// сбрасывает idle-таймер (LLM стримит = жив).
|
||||||
var mu sync.Mutex
|
var mu sync.Mutex
|
||||||
@@ -102,13 +111,18 @@ func (r *Runner) awaitVerdict(ctx context.Context, c *Client, sid, agent string)
|
|||||||
lastProgress := time.Now()
|
lastProgress := time.Now()
|
||||||
launch := 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 {
|
for {
|
||||||
if ctx.Err() != nil {
|
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)
|
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)
|
count, _ := c.textCount(ctx, sid)
|
||||||
@@ -119,80 +133,39 @@ func (r *Runner) awaitVerdict(ctx context.Context, c *Client, sid, agent string)
|
|||||||
}
|
}
|
||||||
mu.Unlock()
|
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()
|
now := time.Now()
|
||||||
if now.Sub(lastProgress) > r.IdleTimeout {
|
if now.Sub(lastProgress) > r.IdleTimeout {
|
||||||
r.logf("opencode(%s) idle %.0fs — abort", agent, r.IdleTimeout.Seconds())
|
r.logf("opencode(%s) idle %.0fs — abort", agent, r.IdleTimeout.Seconds())
|
||||||
if err := c.Abort(ctx, sid); err != nil {
|
return abortAnd(-1, "idle")
|
||||||
r.logf("opencode(%s) abort idle: %v", agent, err)
|
|
||||||
}
|
|
||||||
return &Result{RC: -1, Stdout: "", SessionID: sid}, nil
|
|
||||||
}
|
}
|
||||||
// hard — общий бюджет от старта запуска.
|
// hard — общий бюджет от старта запуска.
|
||||||
if now.Sub(launch) > r.HardTimeout {
|
if now.Sub(launch) > r.HardTimeout {
|
||||||
r.logf("opencode(%s) hard timeout %.0fs — abort", agent, r.HardTimeout.Seconds())
|
r.logf("opencode(%s) hard timeout %.0fs — abort", agent, r.HardTimeout.Seconds())
|
||||||
if err := c.Abort(ctx, sid); err != nil {
|
return abortAnd(-1, "hard")
|
||||||
r.logf("opencode(%s) abort hard: %v", agent, err)
|
|
||||||
}
|
|
||||||
return &Result{RC: -1, Stdout: "", SessionID: sid}, nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
select {
|
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 <-time.After(r.PollInterval):
|
||||||
case <-ctx.Done():
|
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)
|
// ResumeDev — запуск dev-агента с resume-fallback. Если resume (sessionID)
|
||||||
// падает с rc!=0 — повторяем ОДИН раз свежей сессией в том же каталоге.
|
// падает с rc!=0 — повторяем ОДИН раз свежей сессией в том же каталоге.
|
||||||
// rc=-1 (обрыв по таймауту) НЕ триггерит fallback.
|
// rc=-1 (обрыв по таймауту) НЕ триггерит fallback.
|
||||||
|
|||||||
@@ -30,10 +30,7 @@ func httptestURL(t *testing.T, f *fakeAPIServer) string {
|
|||||||
func TestRun_Success(t *testing.T) {
|
func TestRun_Success(t *testing.T) {
|
||||||
dir := t.TempDir()
|
dir := t.TempDir()
|
||||||
f := &fakeAPIServer{
|
f := &fakeAPIServer{
|
||||||
messages: []sessionMessage{{Type: "assistant", Content: []struct {
|
verdictParts: []part{{Type: "text", Text: "done"}},
|
||||||
Type string `json:"type"`
|
|
||||||
Text string `json:"text"`
|
|
||||||
}{textPart("done")}}},
|
|
||||||
}
|
}
|
||||||
p, _ := fakePool(t, f, dir)
|
p, _ := fakePool(t, f, dir)
|
||||||
|
|
||||||
@@ -55,8 +52,8 @@ func TestRun_Success(t *testing.T) {
|
|||||||
|
|
||||||
func TestRun_IdleTimeout(t *testing.T) {
|
func TestRun_IdleTimeout(t *testing.T) {
|
||||||
dir := t.TempDir()
|
dir := t.TempDir()
|
||||||
// сервер никогда не отдаёт text → всегда неготов, прогресс не растёт
|
// prompt блокируется (агент «завис»), прогресс не растёт → idle abort
|
||||||
f := &fakeAPIServer{}
|
f := &fakeAPIServer{blockPrompt: true}
|
||||||
p, _ := fakePool(t, f, dir)
|
p, _ := fakePool(t, f, dir)
|
||||||
|
|
||||||
r := &Runner{Pool: p, IdleTimeout: 30 * time.Millisecond,
|
r := &Runner{Pool: p, IdleTimeout: 30 * time.Millisecond,
|
||||||
@@ -72,7 +69,7 @@ func TestRun_IdleTimeout(t *testing.T) {
|
|||||||
|
|
||||||
func TestRun_ContextCancel(t *testing.T) {
|
func TestRun_ContextCancel(t *testing.T) {
|
||||||
dir := t.TempDir()
|
dir := t.TempDir()
|
||||||
f := &fakeAPIServer{}
|
f := &fakeAPIServer{blockPrompt: true}
|
||||||
p, _ := fakePool(t, f, dir)
|
p, _ := fakePool(t, f, dir)
|
||||||
|
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
@@ -102,10 +99,7 @@ func TestRun_ContextCancel(t *testing.T) {
|
|||||||
func TestResumeDev_NoFallbackOnSuccess(t *testing.T) {
|
func TestResumeDev_NoFallbackOnSuccess(t *testing.T) {
|
||||||
dir := t.TempDir()
|
dir := t.TempDir()
|
||||||
f := &fakeAPIServer{
|
f := &fakeAPIServer{
|
||||||
messages: []sessionMessage{{Type: "assistant", Content: []struct {
|
verdictParts: []part{{Type: "text", Text: "ok"}},
|
||||||
Type string `json:"type"`
|
|
||||||
Text string `json:"text"`
|
|
||||||
}{textPart("ok")}}},
|
|
||||||
}
|
}
|
||||||
p, _ := fakePool(t, f, dir)
|
p, _ := fakePool(t, f, dir)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user