diff --git a/internal/loki/llm_client.go b/internal/loki/llm_client.go index 5349d43..6bf6bd8 100644 --- a/internal/loki/llm_client.go +++ b/internal/loki/llm_client.go @@ -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 ): 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 diff --git a/internal/loki/llm_resume_test.go b/internal/loki/llm_resume_test.go new file mode 100644 index 0000000..b9a5b5b --- /dev/null +++ b/internal/loki/llm_resume_test.go @@ -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) + } +} diff --git a/internal/loki/llm_retry_net.go b/internal/loki/llm_retry_net.go index 47c90fc..f1378d9 100644 --- a/internal/loki/llm_retry_net.go +++ b/internal/loki/llm_retry_net.go @@ -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) +} diff --git a/internal/loki/ui/index.html b/internal/loki/ui/index.html index 2368f9c..2e2c621 100644 --- a/internal/loki/ui/index.html +++ b/internal/loki/ui/index.html @@ -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 diff --git a/internal/loki/ui/src/js/09-stream.js b/internal/loki/ui/src/js/09-stream.js index 695e109..43a0ea3 100644 --- a/internal/loki/ui/src/js/09-stream.js +++ b/internal/loki/ui/src/js/09-stream.js @@ -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