mirror of
https://github.com/R0m1k3/Loki.git
synced 2026-10-11 17:26:57 +02:00
Repris d'AJEAN 0.14.0, adapté. AJOUT EN COURS DE RÉPONSE Il fallait arrêter la génération pour glisser une précision (le serveur répondait 409). Un message envoyé pendant une réponse part désormais EN FILE : runChat l'injecte à la prochaine frontière d'étape (après un appel d'outil) — le modèle en tient compte dans la SUITE de sa réponse — ou, si le tour se termine avant, il devient le tour suivant, dans l'ordre. Le bouton envoyer réapparaît à côté de stop dès qu'il y a du texte ; le message s'affiche « en attente » au-dessus de la carte jusqu'à ce que le flux le confirme. Stop abandonne la file (queue_dropped, le client retire ses pastilles). Écarts avec l'amont : - Les messages en attente sont posés AU-DESSUS de la carte, pas dans le fil qui s'écrit encore (ils s'y seraient intercalés entre deux bulles). - Dédoublonnage par identifiant d'envoi (cid) : l'UI réessaie un envoi dont la réponse s'est perdue sur le tunnel. Le 409 « déjà en cours » faisait office de garde ; sans lui, le réessai mettait le message deux fois en file. - Une tâche planifiée qui occupe le modèle garde le refus 409 (sa fin ne dépile rien). - Pas d'injection après un stop : la boucle repassait en tête une fois l'outil interrompu et journalisait le message en file (vu en test). - runChat prend l'injecteur en paramètre variadique : les appels hors chat (tâches, vérification du mode Code, tests) ne changent pas. DEUX FILS FUSIONNÉS Un appareil déconnecté pendant qu'un autre changeait de discussion se réabonnait avec un `from` hérité de l'ancienne : les Seq n'étant pas globaux, la fin de la nouvelle discussion se greffait sur le début de l'ancienne restée à l'écran. caught_up et reset portent maintenant l'id de la discussion ; le client le renvoie (conv_id) et, s'il ne correspond plus, le serveur ordonne un reset avant de rejouer. La garde « from au-delà du dernier Seq » émet elle aussi ce reset (elle rejouait par-dessus l'écran sans le vider). Vérifié dans le navigateur avec un faux moteur : précision injectée après l'outil dans le même tour, message en file devenu tour suivant, stop qui abandonne la file. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
260 lines
8.9 KiB
Go
260 lines
8.9 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, emit func(map[string]any) bool) {
|
|
tail := -1
|
|
if body.Tail != nil {
|
|
tail = *body.Tail
|
|
}
|
|
conv.SubscribeTail(ctx, body.From, tail, body.ConvID, 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
|
|
}
|
|
// 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()
|
|
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)
|
|
}
|