mirror of
https://github.com/R0m1k3/Loki.git
synced 2026-10-11 17:26:57 +02:00
Agent : reprise après coupure, appels parallèles séparés, outils gardés hors 500
Repris d'AJEAN 0.16.4 et 0.17.4.
- Flux coupé en cours de réponse vers une API externe (Wi-Fi, VPN, proxy
qui décroche) : le tour reprend au lieu d'être abandonné, jusqu'à 5 fois
avec attente croissante. Rien d'affiché : la requête est rejouée (le
raisonnement montré est retiré). Du texte affiché : il est rendu au
modèle avec « continue là où tu t'es arrêté », et la suite s'ajoute à
l'écran. Un outil à moitié écrit est clos puis réémis. Jamais pour le
llama-server local : coupé, il a planté. Une coupure après le dernier
chunk (finish_reason reçu) garde la réponse telle quelle.
- Appels d'outils parallèles : chaque morceau est rangé selon son champ
index, plus selon sa position dans le chunk. Un serveur qui envoie un
appel par chunk les fusionnait (noms écrasés, JSON « {…}{…} »).
- Outils coupés et « réponds maintenant » seulement sur un 500 (appel mal
formé). Un 401, 429 ou 502 d'une API externe coupait les outils et
rejouait le tour, masquant la vraie cause.
- Débit de décodage mesuré côté Loki quand le serveur n'envoie pas de
timings ; le débit de lecture, inconnu, n'est plus affiché à 0.
- Windows : la connexion refusée (WSA 10061) est reconnue comme telle par
la reprise réseau (TestConnexionRefuseeEstRejouee échouait sous Windows).
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
1 parent
2db0ecd431
commit
7f8c5d360f
5 files changed
+312
-16
No files matched your search
+129
-13
@@ -28,6 +28,9 @@ type Message struct {
|
||||
}
|
||||
|
||||
type ToolCall struct {
|
||||
// Index : position de l'appel dans un flux (delta OpenAI). Absent des
|
||||
// messages de l'historique (nil → omis).
|
||||
Index *int `json:"index,omitempty"`
|
||||
ID string `json:"id"`
|
||||
Type string `json:"type"`
|
||||
Function ToolCallFunc `json:"function"`
|
||||
@@ -938,6 +941,11 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
// réponse arrive : une boucle d'outils longue retrouve son budget à chaque
|
||||
// itération réussie, un moteur durablement mort finit par rendre la main.
|
||||
netRetries := 0
|
||||
// Reprises après un flux coupé EN COURS de réponse (API externe seulement,
|
||||
// voir après la lecture du flux). Remis à zéro à chaque complétion lue en
|
||||
// entier.
|
||||
const maxStreamRetries = 5
|
||||
streamRetries := 0
|
||||
// Destination des complétions : llama-server local, ou une API OpenAI-compatible
|
||||
// externe si le preset actif en est un (backend_external.go). Résolu UNE FOIS
|
||||
// par tour — une bascule de preset en plein tour est rare, et se rejoue de
|
||||
@@ -1017,6 +1025,7 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
// l'envoyer à api.openai.com serait fuiter un secret chez un tiers en même
|
||||
// temps qu'un 401 garanti. Chaque endpoint porte la sienne.
|
||||
ep.auth(req.Header.Set)
|
||||
tReq := time.Now() // secours des stats quand le serveur n'envoie pas de timings
|
||||
resp, err := http.DefaultClient.Do(req)
|
||||
if err != nil {
|
||||
// Rien n'est encore parti à l'écran : le tour est rejouable tel quel.
|
||||
@@ -1085,8 +1094,10 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
}
|
||||
// Most common 500 here: llama.cpp couldn't parse a malformed tool call
|
||||
// the model emitted. Retry the turn once without tools so it answers
|
||||
// in plain text rather than leaving the chat dead.
|
||||
if !disableTools && len(tools) > 0 {
|
||||
// in plain text rather than leaving the chat dead. Seulement sur 500 :
|
||||
// un 401, 429 ou 502 (API externe, passerelle) n'a rien à voir avec un
|
||||
// appel mal formé, et couper les outils pour ça cassait le tour d'agent.
|
||||
if resp.StatusCode == http.StatusInternalServerError && !disableTools && len(tools) > 0 {
|
||||
disableTools = true
|
||||
// Nudge the model to answer in plain text from what it already
|
||||
// gathered, so it doesn't immediately re-emit a tool call that
|
||||
@@ -1125,8 +1136,36 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
// arrivent sur des chunks séparés ; on émet une copie complète à chaque MAJ
|
||||
// pour que les consommateurs (terminal, web) aient toujours tout.
|
||||
var stats StatsEvent
|
||||
// Serveur sans `timings` (API tierces) : on mesure nous-mêmes le décodage
|
||||
// entre le premier et le dernier token (voir après la boucle).
|
||||
sawTimings := false
|
||||
var tFirst, tLast time.Time
|
||||
usageGen, genChunks := 0, 0
|
||||
lastPreview := "" // last command preview emitted (to stream the typing)
|
||||
lastBodyLines := -1 // lignes déjà diffusées du corps en cours d'écriture
|
||||
// sentAnswer : du texte de réponse ou un outil est déjà parti vers l'UI pour
|
||||
// cette complétion (une reprise le doublerait). sentReasoning : seul du
|
||||
// raisonnement est parti, qu'on sait retirer (DropReasoning). shown : texte
|
||||
// de réponse réellement affiché, rendu au modèle pour qu'il reprenne après
|
||||
// une coupure. typingTool : appel d'outil en cours d'écriture à l'écran.
|
||||
sentAnswer, sentReasoning := false, false
|
||||
var shown strings.Builder
|
||||
var typingTool *ToolUsedEvent
|
||||
scb := func(ev StreamEvent) bool {
|
||||
if ev.Content != "" {
|
||||
sentAnswer = true
|
||||
shown.WriteString(ev.Content)
|
||||
}
|
||||
if ev.ToolUsed != nil {
|
||||
sentAnswer = true
|
||||
t := *ev.ToolUsed
|
||||
typingTool = &t
|
||||
}
|
||||
if ev.Reasoning != "" {
|
||||
sentReasoning = true
|
||||
}
|
||||
return cb(ev)
|
||||
}
|
||||
// Per-completion reasoning-split state (see reasoningOn comment above).
|
||||
sawReasoningField := false
|
||||
thinkOpen := reasoningOn
|
||||
@@ -1156,6 +1195,7 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
// garde de choices, sinon `gen_tokens`/`gen_per_second` (decode) sont jetés
|
||||
// et l'UI retombe à « 0 tok/s » à la fin de la génération.
|
||||
if chunk.Timings != nil {
|
||||
sawTimings = true
|
||||
stats.PromptTokens = chunk.Timings.PromptN
|
||||
stats.PromptPerSecond = chunk.Timings.PromptPerSecond
|
||||
stats.PromptMs = chunk.Timings.PromptMs
|
||||
@@ -1163,17 +1203,28 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
stats.GenPerSecond = chunk.Timings.PredictedPerSec
|
||||
stats.GenMs = chunk.Timings.PredictedMs
|
||||
s := stats
|
||||
cb(StreamEvent{Stats: &s})
|
||||
scb(StreamEvent{Stats: &s})
|
||||
}
|
||||
if chunk.Usage != nil && chunk.Usage.CompletionTokens > 0 {
|
||||
usageGen = chunk.Usage.CompletionTokens
|
||||
}
|
||||
if chunk.Usage != nil && chunk.Usage.PromptTokens > 0 {
|
||||
stats.PromptTokensTotal = chunk.Usage.PromptTokens
|
||||
s := stats
|
||||
cb(StreamEvent{Stats: &s})
|
||||
scb(StreamEvent{Stats: &s})
|
||||
}
|
||||
if len(chunk.Choices) == 0 {
|
||||
continue
|
||||
}
|
||||
ch := chunk.Choices[0]
|
||||
if ch.Delta.Content != "" || ch.Delta.ReasoningContent != "" || len(ch.Delta.ToolCalls) > 0 {
|
||||
now := time.Now()
|
||||
if tFirst.IsZero() {
|
||||
tFirst = now
|
||||
}
|
||||
tLast = now
|
||||
genChunks++
|
||||
}
|
||||
if ch.FinishReason != "" {
|
||||
finishReason = ch.FinishReason
|
||||
}
|
||||
@@ -1192,14 +1243,21 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
if thinkOpen && thinkTail.Len() > 0 {
|
||||
tail := thinkTail.String()
|
||||
thinkTail.Reset()
|
||||
if !cb(StreamEvent{Content: tail}) {
|
||||
if !scb(StreamEvent{Content: tail}) {
|
||||
aborted = true
|
||||
break
|
||||
}
|
||||
}
|
||||
for i, tc := range ch.Delta.ToolCalls {
|
||||
// llama.cpp's stream may omit index; fall back to slot i.
|
||||
// Le champ index dit à quel appel appartient ce morceau. Un serveur
|
||||
// qui envoie un morceau par chunk met l'appel n°2 en position 0 du
|
||||
// chunk : se fier à i fusionnait les appels parallèles (noms
|
||||
// écrasés, JSON collés « {…}{…} »). Sans index (certains builds
|
||||
// llama.cpp), on retombe sur la position dans le chunk.
|
||||
idx := i
|
||||
if tc.Index != nil {
|
||||
idx = *tc.Index
|
||||
}
|
||||
cur, ok := toolCalls[idx]
|
||||
if !ok {
|
||||
cur = &ToolCall{Type: "function"}
|
||||
@@ -1255,7 +1313,7 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
if grew {
|
||||
lastBodyLines = strings.Count(body, "\n")
|
||||
}
|
||||
if !cb(StreamEvent{ToolUsed: &ToolUsedEvent{Name: cur.Function.Name, Label: lastPreview, Body: body, Typing: true}}) {
|
||||
if !scb(StreamEvent{ToolUsed: &ToolUsedEvent{Name: cur.Function.Name, Label: lastPreview, Body: body, Typing: true}}) {
|
||||
aborted = true
|
||||
break
|
||||
}
|
||||
@@ -1267,7 +1325,7 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
// Backend already separates reasoning — trust it, disable our split.
|
||||
sawReasoningField = true
|
||||
thinkOpen = false
|
||||
if !cb(StreamEvent{Reasoning: ch.Delta.ReasoningContent}) {
|
||||
if !scb(StreamEvent{Reasoning: ch.Delta.ReasoningContent}) {
|
||||
aborted = true
|
||||
break
|
||||
}
|
||||
@@ -1275,7 +1333,7 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
if ch.Delta.Content != "" {
|
||||
assistantContent.WriteString(ch.Delta.Content)
|
||||
if !thinkOpen || sawReasoningField {
|
||||
if !cb(StreamEvent{Content: ch.Delta.Content}) {
|
||||
if !scb(StreamEvent{Content: ch.Delta.Content}) {
|
||||
aborted = true
|
||||
break
|
||||
}
|
||||
@@ -1299,11 +1357,11 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
after := strings.TrimLeft(s[i+len(thinkClose):], "\r\n")
|
||||
thinkOpen = false
|
||||
thinkTail.Reset()
|
||||
if reason != "" && !cb(StreamEvent{Reasoning: reason}) {
|
||||
if reason != "" && !scb(StreamEvent{Reasoning: reason}) {
|
||||
aborted = true
|
||||
break
|
||||
}
|
||||
if after != "" && !cb(StreamEvent{Content: after}) {
|
||||
if after != "" && !scb(StreamEvent{Content: after}) {
|
||||
aborted = true
|
||||
break
|
||||
}
|
||||
@@ -1320,7 +1378,7 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
emit := s[:cut]
|
||||
thinkTail.Reset()
|
||||
thinkTail.WriteString(s[cut:])
|
||||
if !cb(StreamEvent{Content: emit}) {
|
||||
if !scb(StreamEvent{Content: emit}) {
|
||||
aborted = true
|
||||
break
|
||||
}
|
||||
@@ -1331,7 +1389,7 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
}
|
||||
// Flush the held-back tail (never part of a </think>): it's answer text.
|
||||
if !aborted && thinkOpen && thinkTail.Len() > 0 {
|
||||
cb(StreamEvent{Content: strings.TrimLeft(thinkTail.String(), "\r\n")})
|
||||
scb(StreamEvent{Content: strings.TrimLeft(thinkTail.String(), "\r\n")})
|
||||
}
|
||||
// ⚠️ Le flux a-t-il fini, ou CASSÉ ? sc.Scan() renvoie false dans les deux
|
||||
// cas, et l'erreur n'était jamais consultée : une lecture coupée en plein
|
||||
@@ -1344,14 +1402,72 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
// arguments tronqués. On refuse donc le tour, en le disant.
|
||||
scanErr := sc.Err()
|
||||
resp.Body.Close()
|
||||
if !sawTimings && !tFirst.IsZero() {
|
||||
// Décodage mesuré ici ; le 1er token tombe dans la lecture du prompt.
|
||||
// Pas de débit de lecture : le serveur a pu réutiliser une partie du
|
||||
// prompt en cache, et total/délai donnerait des chiffres fantaisistes.
|
||||
gen := usageGen
|
||||
if gen == 0 {
|
||||
gen = genChunks // sans usage : un chunk ≈ un token
|
||||
}
|
||||
stats.PromptMs = float64(tFirst.Sub(tReq).Milliseconds())
|
||||
stats.GenTokens = gen
|
||||
stats.GenMs = float64(tLast.Sub(tFirst).Milliseconds())
|
||||
if sec := tLast.Sub(tFirst).Seconds(); sec > 0 && gen > 1 {
|
||||
stats.GenPerSecond = float64(gen-1) / sec
|
||||
}
|
||||
s := stats
|
||||
cb(StreamEvent{Stats: &s})
|
||||
}
|
||||
if aborted {
|
||||
return extra, nil
|
||||
}
|
||||
// Coupure APRÈS le dernier chunk (finish_reason reçu) : la réponse est
|
||||
// complète, seule la fermeture a raté. On la garde telle quelle.
|
||||
if scanErr != nil && finishReason != "" {
|
||||
scanErr = nil
|
||||
}
|
||||
// Flux coupé vers une API DISTANTE (Wi-Fi, VPN, proxy qui décroche) : on
|
||||
// relance au lieu d'abandonner le tour (AJEAN 0.17.4). Jamais pour le
|
||||
// llama-server local : coupé, il a planté, et rejouer ne ferait que
|
||||
// retarder le message d'erreur.
|
||||
// - Rien de la réponse n'est encore affiché : on rejoue la même requête
|
||||
// (le raisonnement déjà montré est retiré, il va être régénéré).
|
||||
// - Du texte est déjà affiché : on le rend au modèle comme début de sa
|
||||
// réponse, suivi d'un « connexion coupée, continue », et il reprend où
|
||||
// il s'était arrêté. Ces deux messages ne vont que dans la vue modèle
|
||||
// de ce tour : la réponse persistée est le texte complet, reconstitué
|
||||
// par l'appelant à partir du flux. Un appel d'outil à moitié écrit est
|
||||
// jeté : le modèle le réémettra.
|
||||
if scanErr != nil && ctx.Err() == nil && !errors.Is(scanErr, bufio.ErrTooLong) &&
|
||||
ep.External && streamRetries < maxStreamRetries {
|
||||
streamRetries++
|
||||
if sentReasoning && !sentAnswer {
|
||||
cb(StreamEvent{DropReasoning: true})
|
||||
}
|
||||
// Clôt la bulle de l'outil à moitié écrit : il va être réémis, sinon
|
||||
// elle resterait « en cours d'écriture » pour toujours.
|
||||
if typingTool != nil {
|
||||
cb(StreamEvent{ToolUsed: &ToolUsedEvent{Name: typingTool.Name, Label: typingTool.Label, Done: true, Result: "(appel interrompu par une coupure réseau, relancé)"}})
|
||||
}
|
||||
logStreamRetry(streamRetries, maxStreamRetries, scanErr)
|
||||
if llmNetBackoff(ctx, streamRetries) == nil {
|
||||
if partial := shown.String(); strings.TrimSpace(partial) != "" {
|
||||
messages = append(messages,
|
||||
Message{Role: "assistant", Content: partial},
|
||||
Message{Role: "user", Content: "The connection was cut during your answer. Continue exactly where you stopped, without repeating what you already wrote."})
|
||||
}
|
||||
continue
|
||||
}
|
||||
}
|
||||
if scanErr != nil && ctx.Err() == nil {
|
||||
err := streamCutError(scanErr)
|
||||
cb(StreamEvent{Err: err})
|
||||
return extra, err
|
||||
}
|
||||
// Complétion lue en entier : le budget de reprises vaut par coupure
|
||||
// rapprochée, pas pour tout un long tour d'agent.
|
||||
streamRetries = 0
|
||||
|
||||
// Treat any accumulated tool calls as a tool turn even if the backend set
|
||||
// finish_reason to "stop" instead of "tool_calls" (some llama.cpp builds
|
||||
|
||||
@@ -0,0 +1,167 @@
|
||||
package loki
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"net/url"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Flux coupé vers une API EXTERNE après un début de réponse : le tour reprend,
|
||||
// le modèle reçoit le texte déjà affiché suivi d'un « continue », et la suite
|
||||
// s'ajoute à l'écran sans doublon ni erreur.
|
||||
func TestRepriseApresCoupureAPIExterne(t *testing.T) {
|
||||
testHome(t)
|
||||
var n int32
|
||||
var second []Message
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "text/event-stream")
|
||||
if atomic.AddInt32(&n, 1) == 1 {
|
||||
w.WriteHeader(200)
|
||||
_, _ = w.Write([]byte(sseChunk("Bonjour ")))
|
||||
w.(http.Flusher).Flush()
|
||||
conn, _, err := w.(http.Hijacker).Hijack()
|
||||
if err == nil {
|
||||
_ = conn.Close()
|
||||
}
|
||||
return
|
||||
}
|
||||
var body struct {
|
||||
Messages []Message `json:"messages"`
|
||||
}
|
||||
_ = json.NewDecoder(r.Body).Decode(&body)
|
||||
second = body.Messages
|
||||
w.WriteHeader(200)
|
||||
_, _ = w.Write([]byte(sseChunk("le monde")))
|
||||
_, _ = w.Write([]byte(`data: {"choices":[{"delta":{},"finish_reason":"stop"}]}` + "\n\ndata: [DONE]\n\n"))
|
||||
}))
|
||||
t.Cleanup(srv.Close)
|
||||
if err := WriteConfig(parseEnv(externalPresetContent(srv.URL, "m", "", "", false))); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
var content strings.Builder
|
||||
var gotErr error
|
||||
_, err := runChat(context.Background(), []Message{{Role: "user", Content: "salut"}}, 0.7, Caps{}, func(ev StreamEvent) bool {
|
||||
if ev.Err != nil {
|
||||
gotErr = ev.Err
|
||||
}
|
||||
content.WriteString(ev.Content)
|
||||
return true
|
||||
})
|
||||
if err != nil || gotErr != nil {
|
||||
t.Fatalf("coupure externe non reprise : err=%v, ev=%v", err, gotErr)
|
||||
}
|
||||
if got := content.String(); got != "Bonjour le monde" {
|
||||
t.Fatalf("texte affiché = %q", got)
|
||||
}
|
||||
if len(second) < 3 || second[len(second)-2].Role != "assistant" || second[len(second)-2].Content != "Bonjour " ||
|
||||
second[len(second)-1].Role != "user" {
|
||||
t.Fatalf("la reprise n'a pas rendu le début de réponse au modèle : %+v", second)
|
||||
}
|
||||
}
|
||||
|
||||
// Deux appels d'outils parallèles livrés chacun en position 0 de son chunk : le
|
||||
// champ index les sépare. Avant, ils fusionnaient en un seul appel aux
|
||||
// arguments collés « {…}{…} ».
|
||||
func TestAppelsParallelesSuiventLIndex(t *testing.T) {
|
||||
withWorkspace(t)
|
||||
var n int32
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "text/event-stream")
|
||||
w.WriteHeader(200)
|
||||
if atomic.AddInt32(&n, 1) == 1 {
|
||||
_, _ = w.Write([]byte(`data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"a","type":"function","function":{"name":"glob","arguments":"{\"pattern\":\"*.go\"}"}}]}}]}` + "\n\n"))
|
||||
_, _ = w.Write([]byte(`data: {"choices":[{"delta":{"tool_calls":[{"index":1,"id":"b","type":"function","function":{"name":"grep","arguments":"{\"pattern\":\"rien\"}"}}]}}]}` + "\n\n"))
|
||||
_, _ = w.Write([]byte(`data: {"choices":[{"delta":{},"finish_reason":"tool_calls"}]}` + "\n\n"))
|
||||
} else {
|
||||
_, _ = w.Write([]byte(sseChunk("fini")))
|
||||
_, _ = w.Write([]byte(`data: {"choices":[{"delta":{},"finish_reason":"stop"}]}` + "\n\n"))
|
||||
}
|
||||
_, _ = w.Write([]byte("data: [DONE]\n\n"))
|
||||
}))
|
||||
t.Cleanup(srv.Close)
|
||||
u, _ := url.Parse(srv.URL)
|
||||
if err := SetConfigKey("PORT", u.Port()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
var done []string
|
||||
extra, err := runChat(context.Background(), []Message{{Role: "user", Content: "cherche"}}, 0.7, Caps{Agent: true, Code: true}, func(ev StreamEvent) bool {
|
||||
if ev.ToolUsed != nil && ev.ToolUsed.Done {
|
||||
done = append(done, ev.ToolUsed.Name)
|
||||
}
|
||||
return true
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(done) != 2 || done[0] != "glob" || done[1] != "grep" {
|
||||
t.Fatalf("appels exécutés = %v, attendu [glob grep]", done)
|
||||
}
|
||||
if len(extra) == 0 || len(extra[0].ToolCalls) != 2 {
|
||||
t.Fatalf("message assistant : %+v", extra)
|
||||
}
|
||||
}
|
||||
|
||||
// Serveur sans `timings` (API tierce) : le débit de décodage est mesuré côté
|
||||
// Loki au lieu de rester à zéro.
|
||||
func TestDebitMesureSansTimings(t *testing.T) {
|
||||
testHome(t)
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "text/event-stream")
|
||||
w.WriteHeader(200)
|
||||
for _, s := range []string{"un ", "deux ", "trois"} {
|
||||
_, _ = w.Write([]byte(sseChunk(s)))
|
||||
w.(http.Flusher).Flush()
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
}
|
||||
_, _ = w.Write([]byte(`data: {"choices":[{"delta":{},"finish_reason":"stop"}]}` + "\n\n"))
|
||||
_, _ = w.Write([]byte(`data: {"choices":[],"usage":{"prompt_tokens":12,"completion_tokens":3,"total_tokens":15}}` + "\n\ndata: [DONE]\n\n"))
|
||||
}))
|
||||
t.Cleanup(srv.Close)
|
||||
if err := WriteConfig(parseEnv(externalPresetContent(srv.URL, "m", "", "", false))); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var last *StatsEvent
|
||||
if _, err := runChat(context.Background(), []Message{{Role: "user", Content: "compte"}}, 0.7, Caps{}, func(ev StreamEvent) bool {
|
||||
if ev.Stats != nil {
|
||||
last = ev.Stats
|
||||
}
|
||||
return true
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if last == nil || last.GenTokens != 3 || last.GenPerSecond <= 0 {
|
||||
t.Fatalf("stats sans timings = %+v", last)
|
||||
}
|
||||
}
|
||||
|
||||
// Un refus qui n'a rien d'un appel d'outil mal formé (401 d'une API externe)
|
||||
// remonte tel quel : avant, n'importe quel statut coupait les outils et
|
||||
// rejouait le tour, masquant la vraie cause derrière une réponse sans outils.
|
||||
func TestRefusAPIHors500NeCoupePasLesOutils(t *testing.T) {
|
||||
testHome(t)
|
||||
var n int32
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
atomic.AddInt32(&n, 1)
|
||||
w.WriteHeader(401)
|
||||
_, _ = w.Write([]byte(`{"error":{"message":"Incorrect API key provided"}}`))
|
||||
}))
|
||||
t.Cleanup(srv.Close)
|
||||
if err := WriteConfig(parseEnv(externalPresetContent(srv.URL, "m", "", "", false))); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, err := runChat(context.Background(), []Message{{Role: "user", Content: "salut"}}, 0.7, Caps{Agent: true}, func(StreamEvent) bool { return true })
|
||||
if err == nil || !strings.Contains(err.Error(), "401") {
|
||||
t.Fatalf("erreur = %v, attendu le 401 de l'API", err)
|
||||
}
|
||||
if got := atomic.LoadInt32(&n); got != 1 {
|
||||
t.Fatalf("%d requêtes : un 401 a été rejoué sans outils", got)
|
||||
}
|
||||
}
|
||||
@@ -43,6 +43,11 @@ const (
|
||||
llmNetMaxWait = 5 * time.Second
|
||||
)
|
||||
|
||||
// wsaECONNREFUSED : la connexion refusée telle que Windows la rend (WSA 10061).
|
||||
// syscall.ECONNREFUSED n'y est qu'une valeur inventée par Go, jamais renvoyée
|
||||
// par le système ; ailleurs, 10061 n'est le numéro d'aucune erreur.
|
||||
const wsaECONNREFUSED = syscall.Errno(10061)
|
||||
|
||||
// llmRetryableErr : cette erreur de transport vaut-elle une nouvelle tentative ?
|
||||
//
|
||||
// Un ctx déjà fini n'en vaut JAMAIS une : c'est le /stop de l'utilisateur ou la
|
||||
@@ -57,7 +62,7 @@ func llmRetryableErr(ctx context.Context, err error) bool {
|
||||
return false
|
||||
}
|
||||
switch {
|
||||
case errors.Is(err, syscall.ECONNREFUSED): // moteur en cours de démarrage
|
||||
case errors.Is(err, syscall.ECONNREFUSED), errors.Is(err, wsaECONNREFUSED): // moteur en cours de démarrage
|
||||
return true
|
||||
case errors.Is(err, syscall.ECONNRESET), errors.Is(err, io.EOF), errors.Is(err, io.ErrUnexpectedEOF):
|
||||
return true
|
||||
@@ -122,3 +127,9 @@ func llmNetBackoff(ctx context.Context, n int) error {
|
||||
func logLLMRetry(attempt int, cause string) {
|
||||
fmt.Fprintf(os.Stderr, "[llm] %s — nouvelle tentative %d/%d\n", cause, attempt, llmNetRetries)
|
||||
}
|
||||
|
||||
// logStreamRetry trace la reprise d'un flux coupé EN COURS de réponse (API
|
||||
// externe) : budget distinct de celui d'avant le flux, voir runChat.
|
||||
func logStreamRetry(attempt, max int, cause error) {
|
||||
fmt.Fprintf(os.Stderr, "[llm] flux coupé en cours de réponse (%v) — reprise %d/%d\n", cause, attempt, max)
|
||||
}
|
||||
@@ -7211,7 +7211,8 @@ function renderStats(el, s){
|
||||
if(!el||!s) return;
|
||||
const parts=[];
|
||||
const pt=s.prompt_tokens||s.prompt_tokens_total;
|
||||
if(pt) parts.push('prefill '+pt+' tok · '+(s.prompt_per_second||0).toFixed(0)+' tok/s');
|
||||
// Débit de lecture inconnu (API tierce sans timings) : on n'affiche que la taille.
|
||||
if(pt) parts.push('prefill '+pt+' tok'+(s.prompt_per_second ? ' · '+s.prompt_per_second.toFixed(0)+' tok/s' : ''));
|
||||
if(s.gen_tokens) parts.push('decode '+s.gen_tokens+' tok · '+(s.gen_per_second||0).toFixed(1)+' tok/s');
|
||||
if(!parts.length) return;
|
||||
// Réponse de l'assistant : ligne de mesures dédiée sous le texte (son étiquette
|
||||
|
||||
@@ -237,7 +237,8 @@ function renderStats(el, s){
|
||||
if(!el||!s) return;
|
||||
const parts=[];
|
||||
const pt=s.prompt_tokens||s.prompt_tokens_total;
|
||||
if(pt) parts.push('prefill '+pt+' tok · '+(s.prompt_per_second||0).toFixed(0)+' tok/s');
|
||||
// Débit de lecture inconnu (API tierce sans timings) : on n'affiche que la taille.
|
||||
if(pt) parts.push('prefill '+pt+' tok'+(s.prompt_per_second ? ' · '+s.prompt_per_second.toFixed(0)+' tok/s' : ''));
|
||||
if(s.gen_tokens) parts.push('decode '+s.gen_tokens+' tok · '+(s.gen_per_second||0).toFixed(1)+' tok/s');
|
||||
if(!parts.length) return;
|
||||
// Réponse de l'assistant : ligne de mesures dédiée sous le texte (son étiquette
|
||||
|
||||
Reference in new issue
Block a user