mirror of
https://github.com/R0m1k3/Loki.git
synced 2026-10-12 01:37:06 +02:00
Deux manques dans la présentation d'un tour : 1. Rien ne disait sur QUELLE discussion un agent était en train de travailler. La ligne concernée porte maintenant un anneau qui tourne, à gauche du titre — et l'en-tête aussi, puisque la barre latérale est escamotée sur téléphone. L'état vient de deux sources : le flux SSE (instantané) et /api/conversations, qui expose désormais `busy` pour que la page ouverte en cours de génération s'anime sans attendre le rejeu du journal. L'indicateur est posé et retiré SANS redessiner la liste : un redessin relancerait l'animation à chaque événement. 2. Le pied des réponses ne montrait que la vitesse du moteur, jamais le temps que le tour avait pris. Il porte maintenant « travail 1 min 04 s » — de la question envoyée à la fin du tour, raisonnement, appels d'outils et attentes compris. La durée avance à la seconde pendant le tour et se fige au turn_done ; elle apparaît aussi sur l'indicateur « … », là où il n'y a pas encore de réponse sous laquelle écrire. La bulle de raisonnement gagne au passage sa propre durée. La mesure est prise sur les HORODATAGES SERVEUR (événements `user` et `turn_done`) : elle est donc juste en direct comme au rejeu, où tout arrive d'un bloc côté client. Le compteur vivant se recale sur l'écart entre l'horloge du navigateur et celle du serveur, réévalué à chaque événement reçu en direct — un téléphone n'est pas à la même heure que la machine. « Masquer la vitesse de génération » ne masque plus que la vitesse : la durée n'est pas une mesure de moteur, c'est ce que la réponse a coûté. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01J5UndZ9DedRXPuAoRXbmDb
225 lines
7.6 KiB
Go
225 lines
7.6 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
|
|
}
|
|
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, emit func(map[string]any) bool) {
|
|
conv.Subscribe(ctx, body.From, emit)
|
|
}
|
|
|
|
// 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
|
|
}
|
|
if err := conv.StartTurn(body.Message, files, capsFromBody(body), body.Temperature); err != nil {
|
|
// 409 = occupé (génération en cours) ; 503 = modèle pas prêt.
|
|
code := 503
|
|
if err == ErrBusy {
|
|
code = 409
|
|
}
|
|
sendJSON(w, code, map[string]any{"ok": false, "error": err.Error()})
|
|
return
|
|
}
|
|
sendJSON(w, 200, map[string]any{"ok": true})
|
|
}
|
|
|
|
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
|
|
}
|
|
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()
|
|
emit := func(obj map[string]any) bool {
|
|
b, _ := json.Marshal(map[string]any{"choices": []any{map[string]any{"delta": obj}}})
|
|
mu.Lock()
|
|
defer mu.Unlock()
|
|
if _, err := w.Write([]byte("data: " + string(b) + "\n\n")); err != nil {
|
|
return false
|
|
}
|
|
if flusher != nil {
|
|
flusher.Flush()
|
|
}
|
|
return true
|
|
}
|
|
runChatStream(r.Context(), body, emit)
|
|
}
|