Files
Loki/internal/loki/web_chat.go
T
MichaelandClaude Opus 5.5 f1f6cbe4ea Flux : l'écriture d'un gros fichier ne coûte plus un temps quadratique côté Go
Pendant qu'un modèle tapait un write de 120 Ko, chaque morceau faisait relire
et redécoder tout le JSON des arguments, puis republiait le corps ENTIER à
chaque ligne : 250 Mo de flux SSE, autant de tas vivant dans le journal, et
près de 8 s de CPU volées au décodage quand des experts tournent sur le
processeur. Même travers pour la détection des appels d'outils écrits en
texte, qui repassait toutes les regex sur toute la réponse à chaque jeton.
Rien de ce que voit le modèle ne change : les arguments exécutés restent
entiers, la passe complète de fin de tour aussi.

- argPreview : lecture incrémentale de l'argument affiché et du corps, même
  règle que previewArgDone (barre oblique coupée gardée en attente) ;
  arguments accumulés dans un strings.Builder, cur.Function.Arguments recalé
  sur lui à chaque morceau.
- Événement de frappe : seules les 40 dernières lignes (4 Kio au plus) du
  corps, avec body_lines (vrai « +N ») et body_tail ; l'UI affiche « … » et
  garde le repli sur le décompte local.
- textualToolCallFrom : ne relit que la fin, depuis un vrai début de ligne
  avant la fin du balayage précédent, en reculant sur les suites que les
  motifs embarqués peuvent traverser. Surcharge retry_patterns : balayage
  complet, comme avant.
- Direct SSE : les événements trouvés à un réveil partent en une écriture
  (lots de 128 Kio / 256 événements au plus), un sceau par événement en E2E.
- Tests : fuzz d'argPreview contre previewArgDone, équivalence fenêtre /
  balayage complet au morceau près, corps borné de bout en bout, ordre des
  lots ; bancs d'essai.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-04 01:48:49 +02:00

317 lines
11 KiB
Go

// web_chat.go — endpoints de chat du serveur web local : envoi/stop/reset,
// flux SSE d'abonnement à la conversation serveur (voir chat_conversation.go).
package loki
import (
"context"
"encoding/json"
"net/http"
"strings"
"sync"
"time"
)
// capsFromBody dérive les capacités d'un tour à partir de la configuration
// machine et des surcharges portées par la requête.
//
// ⚠️ Une surcharge ne peut que RESTREINDRE. Elle pouvait auparavant rallumer le
// mode agent : un simple {"agent":true} redonnait bash, write, edit et les
// outils MCP alors que l'interrupteur de la machine était sur OFF. Comme l'API
// n'est pas protégée par défaut et écoute sur 0.0.0.0, l'interrupteur ne
// garantissait donc rien. Il redevient une vraie fermeture : ce qui est éteint
// sur la machine ne peut pas être rallumé par un client.
func capsFromBody(body chatReq) Caps {
caps := globalCaps()
switch {
case body.Agent != nil:
caps.Agent = caps.Agent && *body.Agent
case body.Tools != nil || body.Skills != nil:
want := (body.Tools != nil && *body.Tools) || (body.Skills != nil && *body.Skills)
caps.Agent = caps.Agent && want
}
if body.Internet != nil {
caps.Internet = caps.Internet && *body.Internet
}
// Les outils dépendent du mode agent : agent coupé, tout est coupé.
if !caps.Agent {
caps.Internet = false
}
// Mode code : réglage PAR DISCUSSION (sélecteur Chat/Code du composeur),
// pas une surcharge de requête — même logique de fermeture que le reste :
// sans agent, pas de mode code.
caps.Code = caps.Agent && convCodeMode()
return caps
}
// sseHeartbeat garde la réponse SSE active en écrivant un commentaire (`: ping`,
// ignoré par le parseur côté navigateur, aucun contenu donc rien à chiffrer)
// toutes les ~15 s. Sans ça, un long silence (exécution d'outil en mode agent,
// gros prefill) laisse la réponse inactive et un proxy intermédiaire (Cloudflare,
// ~100 s) la coupe → le fetch navigateur échoue (« Load failed »). Retourne un
// mutex à partager avec l'émetteur (writes concurrents sur le même w) et une
// fonction d'arrêt à différer.
func sseHeartbeat(w http.ResponseWriter, flusher http.Flusher) (*sync.Mutex, func()) {
mu := &sync.Mutex{}
done := make(chan struct{})
go func() {
// 4 s (et non 15) : borne le temps qu'un dernier bout de flux peut rester
// coincé dans un buffer proxy (Cloudflare) faute d'octets pour le pousser.
t := time.NewTicker(4 * time.Second)
defer t.Stop()
for {
select {
case <-done:
return
case <-t.C:
mu.Lock()
_, err := w.Write([]byte(": ping\n\n"))
if flusher != nil {
flusher.Flush()
}
mu.Unlock()
if err != nil {
return
}
}
}
}()
return mu, func() { close(done) }
}
// runChatStream est désormais un pur ABONNÉ au journal de la conversation serveur :
// il rejoue Log[body.From:] puis suit le direct, jusqu'à ce que la connexion (ctx)
// se ferme. La GÉNÉRATION est lancée séparément par /api/chat/send dans une
// goroutine détachée — fermer le navigateur n'arrête donc plus rien. Partagé par
// handleChat (clair) et handleE2EChat (chiffré).
func runChatStream(ctx context.Context, body chatReq, out *sseStream) {
tail := -1
if body.Tail != nil {
tail = *body.Tail
}
conv.subscribeSink(ctx, body.From, tail, body.ConvID, tailSink{emit: out.emit, queue: out.queue, flush: out.flush})
}
// sseBatchBytes / sseBatchEvents : taille d'un envoi groupé. Un client très en
// retard peut trouver des milliers d'événements en attente : on écrit par
// tranches au lieu de tout empiler en mémoire.
const (
sseBatchBytes = 128 << 10
sseBatchEvents = 256
)
// sseStream : écriture des trames `data:` d'un flux d'abonnement. frame met un
// événement en forme (en clair, ou scellé pour le relais E2E — un sceau par
// événement, le client déchiffre chaque ligne `data:` séparément). mu est
// celui du battement de cœur : jamais deux écritures mêlées sur w.
type sseStream struct {
w http.ResponseWriter
flusher http.Flusher
mu *sync.Mutex
frame func(map[string]any) ([]byte, bool)
buf []byte
n int
}
// emit : un événement, écrit et poussé tout de suite (après ce qui attend).
func (s *sseStream) emit(obj map[string]any) bool {
return s.queue(obj) && s.flush()
}
// queue : un événement de plus dans le tampon, écrit d'office s'il déborde.
func (s *sseStream) queue(obj map[string]any) bool {
b, ok := s.frame(obj)
if !ok {
return false
}
s.buf = append(s.buf, b...)
s.n++
if len(s.buf) >= sseBatchBytes || s.n >= sseBatchEvents {
return s.flush()
}
return true
}
// flush : une seule écriture et un seul flush pour tout ce qui attend.
func (s *sseStream) flush() bool {
if len(s.buf) == 0 {
return true
}
s.mu.Lock()
defer s.mu.Unlock()
_, err := s.w.Write(s.buf)
// Un gros événement (diff final) ne doit pas garder son tampon à vie.
if cap(s.buf) > 4*sseBatchBytes {
s.buf = nil
} else {
s.buf = s.buf[:0]
}
s.n = 0
if err != nil {
return false
}
if s.flusher != nil {
s.flusher.Flush()
}
return true
}
// handleChatSend ajoute un message et lance la génération en arrière-plan. Réponse
// req/resp (les événements arrivent par le flux d'abonnement). Passe par le proxy
// tunnel /api/e2e/req pour app.ajean.link — aucun code E2E spécifique requis.
func handleChatSend(w http.ResponseWriter, r *http.Request) {
var body chatReq
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
sendJSON(w, 400, map[string]any{"ok": false, "error": err.Error()})
return
}
// Un envoi SANS texte mais AVEC pièce jointe est légitime (« tiens, regarde »).
files := attachFiles(body.Files)
if strings.TrimSpace(body.Message) == "" && len(files) == 0 {
sendJSON(w, 400, map[string]any{"ok": false, "error": "message vide"})
return
}
rememberUserTZ(body.TZ)
// Génération en cours : au lieu de refuser, on MET EN FILE (AJEAN 0.14.0).
// Le message est injecté dans la réponse en cours à la prochaine frontière
// d'étape, ou traité comme tour suivant. queued=true le signale au client.
queued, err := conv.EnqueueOrStart(body.ClientID, body.Message, files, capsFromBody(body), body.Temperature)
if err == ErrDupSend {
// Réessai d'un envoi déjà accepté (réponse perdue en route) : succès.
sendJSON(w, 200, map[string]any{"ok": true, "dup": true})
return
}
if err != nil {
// 409 = une tâche planifiée occupe le modèle ; 503 = modèle pas prêt.
code := 503
if err == ErrBusy {
code = 409
err = conv.busyReason() // dit LAQUELLE des deux occupe le modèle
}
sendJSON(w, code, map[string]any{"ok": false, "error": err.Error()})
return
}
sendJSON(w, 200, map[string]any{"ok": true, "queued": queued})
}
// handleToolResult renvoie le résultat COMPLET d'un outil dont le flux n'a
// transporté qu'un aperçu (bouton « voir plus », voir tool_results.go).
func handleToolResult(w http.ResponseWriter, r *http.Request) {
id := strings.TrimSpace(r.URL.Query().Get("id"))
if id == "" {
sendJSON(w, 400, map[string]any{"error": "id manquant"})
return
}
s, ok := loadToolResult(id)
if !ok {
sendJSON(w, 404, map[string]any{"error": "résultat non disponible"})
return
}
sendJSON(w, 200, map[string]any{"result": s})
}
func handleChatStop(w http.ResponseWriter, r *http.Request) {
conv.Stop()
sendJSON(w, 200, map[string]any{"ok": true})
}
func handleChatReset(w http.ResponseWriter, r *http.Request) {
conv.Reset()
sendJSON(w, 200, map[string]any{"ok": true})
}
// --- Discussions (plusieurs fils) ------------------------------------------
// Le fil actif est PARTAGÉ par tous les appareils, comme la conversation unique
// de l'amont : basculer sur le téléphone bascule aussi le PC (l'epoch fait
// rejouer le nouveau fil à tous les abonnés SSE).
func handleConvList(w http.ResponseWriter, r *http.Request) {
list, active := ConvList()
if list == nil {
list = []convMeta{}
}
// `busy` = un tour est en cours SUR LA DISCUSSION ACTIVE (il n'y a qu'une
// génération à la fois, sur le fil ouvert). L'UI s'en sert pour animer la
// ligne concernée dès le chargement, avant même que le flux SSE ait rejoué
// le journal — sinon la page reste muette pendant qu'un agent travaille.
sendJSON(w, 200, map[string]any{"conversations": list, "active": active, "busy": conv.isGenerating()})
}
func handleConvNew(w http.ResponseWriter, r *http.Request) {
sendJSON(w, 200, map[string]any{"ok": true, "id": convNew()})
}
func handleConvSwitch(w http.ResponseWriter, r *http.Request) {
var req struct {
ID string `json:"id"`
}
_ = json.NewDecoder(r.Body).Decode(&req)
if err := convSwitch(req.ID); err != nil {
sendJSON(w, 400, map[string]any{"ok": false, "error": err.Error()})
return
}
sendJSON(w, 200, map[string]any{"ok": true})
}
func handleConvRename(w http.ResponseWriter, r *http.Request) {
var req struct {
ID string `json:"id"`
Title string `json:"title"`
}
_ = json.NewDecoder(r.Body).Decode(&req)
if err := convRename(req.ID, req.Title); err != nil {
sendJSON(w, 400, map[string]any{"ok": false, "error": err.Error()})
return
}
sendJSON(w, 200, map[string]any{"ok": true})
}
func handleConvDelete(w http.ResponseWriter, r *http.Request) {
var req struct {
ID string `json:"id"`
}
_ = json.NewDecoder(r.Body).Decode(&req)
if err := convDelete(req.ID); err != nil {
sendJSON(w, 400, map[string]any{"ok": false, "error": err.Error()})
return
}
sendJSON(w, 200, map[string]any{"ok": true})
}
// handleChatCompact lance une compaction manuelle du contexte (bouton UI). La
// progression est diffusée via le flux d'abonnement (compacting/compacted).
func handleChatCompact(w http.ResponseWriter, r *http.Request) {
if err := conv.CompactNow(); err != nil {
code := 503
if err == ErrBusy {
code = 409
err = conv.busyReason()
}
sendJSON(w, code, map[string]any{"ok": false, "error": err.Error()})
return
}
sendJSON(w, 200, map[string]any{"ok": true})
}
func handleChatState(w http.ResponseWriter, r *http.Request) {
sendJSON(w, 200, conv.state())
}
func handleChat(w http.ResponseWriter, r *http.Request) {
var body chatReq
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
http.Error(w, err.Error(), 400)
return
}
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache, no-transform")
w.Header().Set("X-Accel-Buffering", "no")
flusher, _ := w.(http.Flusher)
mu, stop := sseHeartbeat(w, flusher)
defer stop()
frame := func(obj map[string]any) ([]byte, bool) {
b, _ := json.Marshal(map[string]any{"choices": []any{map[string]any{"delta": obj}}})
return []byte("data: " + string(b) + "\n\n"), true
}
runChatStream(r.Context(), body, &sseStream{w: w, flusher: flusher, mu: mu, frame: frame})
}