mirror of
https://github.com/R0m1k3/Loki.git
synced 2026-10-11 17:26:57 +02:00
Le flux de reponse coupe ne s'arrete plus en silence (issue #19)
sc.Scan() rend false aussi bien a la fin normale d'un flux qu'a la premiere erreur de lecture, et sc.Err() n'etait JAMAIS consulte. Une connexion reinitialisee en plein milieu, ou une ligne plus longue que le tampon, etait donc indiscernable d'une fin propre : le tour s'arretait sans un mot pendant que le journal de llama-server affichait un « stop processing » parfaitement normal. C'est le symptome rapporte — l'agent enchaine quelques commandes puis rend la main sans avoir termine ni commente. Pire : des appels d'outils accumules a moitie auraient ete EXECUTES avec des arguments tronques. On abandonne desormais le tour en nommant la cause, le texte deja recu reste affiche, et un stop volontaire reste silencieux. Trois tests couvrent les trois cas ; celui du flux coupe echoue bien sur l'ancien code. Tampon du scanner porte de 1 a 8 Mio au passage : la cause la plus probable d'une ligne trop longue est un appel d'outil demesure, et jusqu'ici il partait dans le meme silence. Sans rapport, trouve en faisant tourner la suite : renameAside tirait son nom du seul UnixNano, dont la granularite Windows peut rendre deux fois la meme valeur — deux ecartements se disputaient alors le meme nom, le second ecrasant le binaire mis de cote par le premier. Un compteur atomique les departage.
This commit is contained in:
1 parent
4a66b27227
commit
42cdb2cda3
3 files changed
+178
-3
No files matched your search
@@ -517,6 +517,17 @@ func errEngineReset() error {
|
||||
return fmt.Errorf("⚠️ Connexion au moteur (llama-server, port %d) interrompue — il a peut-être redémarré. Réessaie.", LLMPort())
|
||||
}
|
||||
|
||||
// streamCutError explique un flux de complétion coupé en cours de route. Cas à
|
||||
// part de friendlyLLMError : ici la requête avait ABOUTI (200 reçu, tokens déjà
|
||||
// reçus), c'est la lecture qui a lâché. Le dire autrement qu'un « le moteur ne
|
||||
// répond pas » évite d'envoyer l'utilisateur vérifier un moteur qui va bien.
|
||||
func streamCutError(err error) error {
|
||||
if errors.Is(err, bufio.ErrTooLong) {
|
||||
return fmt.Errorf("⚠️ Réponse du moteur illisible : une ligne du flux dépasse la taille maximale (%d Mio). C'est presque toujours un appel d'outil démesuré (écriture d'un très gros fichier). Le tour est abandonné pour ne pas exécuter un appel tronqué.", 8)
|
||||
}
|
||||
return fmt.Errorf("⚠️ Le flux de réponse du moteur (llama-server, port %d) a été coupé en cours de route : %v. La réponse est incomplète et le tour est abandonné — réessaie.", LLMPort(), err)
|
||||
}
|
||||
|
||||
// isNetTimeout : une erreur réseau qui se déclare elle-même comme un délai
|
||||
// dépassé (net.Error.Timeout), quel que soit son libellé.
|
||||
func isNetTimeout(err error) bool {
|
||||
@@ -650,9 +661,12 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
sawReasoningField := false
|
||||
thinkOpen := reasoningOn
|
||||
var thinkTail strings.Builder
|
||||
// scanner with a big buffer — some chunks include large arguments JSON
|
||||
// Scanner à gros tampon : un chunk peut porter un gros JSON d'arguments
|
||||
// (écriture de fichier). 8 Mio et non 1 : au-delà du tampon, le scanner
|
||||
// s'arrête sur « token too long » AU MILIEU du flux, et jusqu'à la 0.8.4
|
||||
// personne ne le voyait (voir sc.Err() plus bas).
|
||||
sc := bufio.NewScanner(resp.Body)
|
||||
sc.Buffer(make([]byte, 0, 64*1024), 1<<20)
|
||||
sc.Buffer(make([]byte, 0, 64*1024), 8<<20)
|
||||
aborted := false
|
||||
for sc.Scan() {
|
||||
line := strings.TrimSpace(sc.Text())
|
||||
@@ -832,10 +846,25 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
|
||||
if !aborted && thinkOpen && thinkTail.Len() > 0 {
|
||||
cb(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
|
||||
// milieu (connexion réinitialisée, ligne plus longue que le tampon) était
|
||||
// donc indiscernable d'une fin normale. Vécu par l'utilisateur : l'agent
|
||||
// enchaîne quelques commandes puis « s'arrête et rend la main sans avoir
|
||||
// terminé, ni même commenté », pendant que le journal de llama-server
|
||||
// affiche un « stop processing » parfaitement normal (issue #19). Pire,
|
||||
// des appels d'outils accumulés à moitié auraient été EXÉCUTÉS avec des
|
||||
// arguments tronqués. On refuse donc le tour, en le disant.
|
||||
scanErr := sc.Err()
|
||||
resp.Body.Close()
|
||||
if aborted {
|
||||
return extra, nil
|
||||
}
|
||||
if scanErr != nil && ctx.Err() == nil {
|
||||
err := streamCutError(scanErr)
|
||||
cb(StreamEvent{Err: err})
|
||||
return extra, err
|
||||
}
|
||||
|
||||
// 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,139 @@
|
||||
package ajean
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"net/url"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// sseServer sert un flux de complétion puis le COUPE brutalement (fermeture de
|
||||
// la connexion sans terminer la réponse), comme le ferait un moteur qui perd sa
|
||||
// connexion en plein milieu. Renvoie le port à mettre dans PORT.
|
||||
func sseCuttingServer(t *testing.T, body string, cut bool) string {
|
||||
t.Helper()
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "text/event-stream")
|
||||
w.WriteHeader(200)
|
||||
_, _ = w.Write([]byte(body))
|
||||
w.(http.Flusher).Flush()
|
||||
if !cut {
|
||||
return
|
||||
}
|
||||
// Coupe la connexion sans fin de réponse propre : côté client, la lecture
|
||||
// du corps échoue (« unexpected EOF ») au lieu de se terminer.
|
||||
hj, ok := w.(http.Hijacker)
|
||||
if !ok {
|
||||
t.Error("serveur de test non détournable")
|
||||
return
|
||||
}
|
||||
conn, _, err := hj.Hijack()
|
||||
if err != nil {
|
||||
t.Error(err)
|
||||
return
|
||||
}
|
||||
_ = conn.Close()
|
||||
}))
|
||||
t.Cleanup(srv.Close)
|
||||
u, err := url.Parse(srv.URL)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return u.Port()
|
||||
}
|
||||
|
||||
// chunk fabrique une ligne SSE de contenu.
|
||||
func sseChunk(text string) string {
|
||||
return `data: {"choices":[{"delta":{"content":"` + text + `"}}]}` + "\n\n"
|
||||
}
|
||||
|
||||
// Le cœur de l'issue #19 : llama-server termine normalement de son côté (« stop
|
||||
// processing » dans son journal), mais la lecture du flux casse. Avant, le tour
|
||||
// s'arrêtait EN SILENCE — l'agent « rendait la main sans avoir terminé, ni même
|
||||
// commenté ». Il doit maintenant remonter une erreur explicite.
|
||||
func TestFluxCoupeRemonteUneErreur(t *testing.T) {
|
||||
testHome(t)
|
||||
// Réponse annoncée sur 3 chunks mais coupée après le premier, sans [DONE].
|
||||
port := sseCuttingServer(t, sseChunk("début de r"), true)
|
||||
if err := SetConfigKey("PORT", port); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
var gotErr error
|
||||
var content strings.Builder
|
||||
_, err := runChat(context.Background(), []Message{{Role: "user", Content: "bonjour"}}, 0.7, Caps{}, func(ev StreamEvent) bool {
|
||||
switch {
|
||||
case ev.Err != nil:
|
||||
gotErr = ev.Err
|
||||
case ev.Content != "":
|
||||
content.WriteString(ev.Content)
|
||||
}
|
||||
return true
|
||||
})
|
||||
if err == nil {
|
||||
t.Fatal("flux coupé : runChat a rendu la main SANS erreur (le tour s'arrêtait en silence)")
|
||||
}
|
||||
if gotErr == nil {
|
||||
t.Fatal("flux coupé : aucune erreur poussée vers l'interface")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "coupé") {
|
||||
t.Errorf("message peu clair pour l'utilisateur : %v", err)
|
||||
}
|
||||
// Le texte déjà reçu n'est pas jeté : on a bien affiché ce qui était arrivé.
|
||||
if !strings.Contains(content.String(), "début de r") {
|
||||
t.Errorf("le texte reçu avant la coupure a été perdu : %q", content.String())
|
||||
}
|
||||
}
|
||||
|
||||
// Contre-épreuve : un flux qui se termine PROPREMENT ne doit évidemment pas
|
||||
// déclencher l'erreur, sans quoi chaque réponse normale se plaindrait.
|
||||
func TestFluxCompletNeRemonteRien(t *testing.T) {
|
||||
testHome(t)
|
||||
port := sseCuttingServer(t, sseChunk("réponse complète")+"data: [DONE]\n\n", false)
|
||||
if err := SetConfigKey("PORT", port); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var gotErr error
|
||||
var content strings.Builder
|
||||
_, err := runChat(context.Background(), []Message{{Role: "user", Content: "bonjour"}}, 0.7, Caps{}, func(ev StreamEvent) bool {
|
||||
if ev.Err != nil {
|
||||
gotErr = ev.Err
|
||||
}
|
||||
if ev.Content != "" {
|
||||
content.WriteString(ev.Content)
|
||||
}
|
||||
return true
|
||||
})
|
||||
if err != nil || gotErr != nil {
|
||||
t.Fatalf("flux normal signalé en erreur : %v / %v", err, gotErr)
|
||||
}
|
||||
if content.String() != "réponse complète" {
|
||||
t.Fatalf("contenu = %q", content.String())
|
||||
}
|
||||
}
|
||||
|
||||
// Un arrêt demandé par l'utilisateur (/stop) coupe aussi la lecture : il ne doit
|
||||
// PAS se déguiser en panne du moteur.
|
||||
func TestFluxAnnuleResteSilencieux(t *testing.T) {
|
||||
testHome(t)
|
||||
port := sseCuttingServer(t, sseChunk("a"), true)
|
||||
if err := SetConfigKey("PORT", port); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
var gotErr error
|
||||
_, _ = runChat(ctx, []Message{{Role: "user", Content: "bonjour"}}, 0.7, Caps{}, func(ev StreamEvent) bool {
|
||||
if ev.Content != "" {
|
||||
cancel() // l'utilisateur clique sur stop dès le premier morceau
|
||||
}
|
||||
if ev.Err != nil {
|
||||
gotErr = ev.Err
|
||||
}
|
||||
return true
|
||||
})
|
||||
if gotErr != nil && strings.Contains(gotErr.Error(), "coupé") {
|
||||
t.Fatalf("un stop volontaire a été présenté comme une panne : %v", gotErr)
|
||||
}
|
||||
}
|
||||
@@ -12,6 +12,7 @@ import (
|
||||
"path/filepath"
|
||||
"runtime"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"golang.org/x/mod/semver"
|
||||
@@ -380,8 +381,14 @@ func replaceBinary(exe, tmp string) error {
|
||||
// permet de remplacer l'image d'un exécutable en cours d'exécution. Reproduit
|
||||
// puis vérifié corrigé sur banc d'essai, deux processus vivants sur les deux
|
||||
// noms : avec un suffixe unique le renommage passe.
|
||||
// asideSeq départage deux écartements rapprochés. UnixNano seul ne suffit PAS :
|
||||
// la granularité de l'horloge Windows peut rendre la même valeur à deux appels
|
||||
// consécutifs, et les deux écartements se disputaient alors le même nom — le
|
||||
// second renommage écrasant le premier, donc le binaire précédent perdu.
|
||||
var asideSeq atomic.Uint64
|
||||
|
||||
func renameAside(path string) (string, error) {
|
||||
aside := fmt.Sprintf("%s.old-%d", path, time.Now().UnixNano())
|
||||
aside := fmt.Sprintf("%s.old-%d-%d", path, time.Now().UnixNano(), asideSeq.Add(1))
|
||||
if err := os.Rename(path, aside); err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user