conversation persistee cote serveur + generation autonome (multi-appareils)

L'historique et la generation vivent desormais sur le serveur jean, plus dans le
localStorage du navigateur :
- conversation.go : store unique (Messages vue modele + Log d'evenements rejouable),
  persistance disque JeanHome/conversation.json, generation dans une goroutine
  detachee (context.Background) -> fermer le navigateur ne coupe plus la reponse
- runChatStream devient un pur abonne (rejoue Log[from:] puis suit le direct) ;
  nouveaux endpoints /api/chat/send|stop|reset|state
- UI : flux d'abonnement permanent + replay du journal (outils/vitesses/raisonnement
  reviennent apres refresh), synchro live inter-appareils, reconnexion auto
- compactage manuel neutralise (auto cote serveur), vitesses rendues depuis les
  stats persistees
This commit is contained in:
nathaninline committed 2026-07-21 17:20:27 +02:00
1 parent 2daea32ade
commit 4514a96585
4 files changed
+605 -265

No files matched your search

+304
View File
@@ -0,0 +1,304 @@
package main
import (
"context"
"encoding/json"
"fmt"
"os"
"path/filepath"
"strings"
"sync"
)
// État de conversation CÔTÉ SERVEUR — une seule conversation partagée par tous
// les appareils. Avant, l'historique vivait dans le localStorage de chaque
// navigateur : refresh = perte des détails (outils/vitesses/raisonnement),
// contexte différent par appareil, et fermer l'onglet coupait la génération
// (liée à r.Context()). Ici l'état est possédé par le serveur, persisté sur
// disque, et la génération tourne dans une goroutine détachée : fermer le
// navigateur ne l'arrête plus, et se reconnecter rejoue tout le fil.
//
// Idée clé : l'UI reconstruit déjà tout l'affichage à partir d'une suite
// d'événements SSE `delta` (content, reasoning_content, tool_used, stats…). On
// JOURNALISE ces événements (Log) horodatés par Seq. Se reconnecter = rejouer
// Log[from:] puis suivre les événements en direct — aucun code de rendu nouveau.
// maxLogEvents plafonne le journal d'AFFICHAGE (pas la vue modèle). Une très
// longue conversation ne fait pas exploser la mémoire/le disque ; la troncature
// ne retire que de vieilles bulles à l'écran, jamais du contexte du modèle
// (celui-ci vit dans Messages, réduit séparément par le compactage).
const maxLogEvents = 5000
// LogEvent = un événement d'affichage rejouable (un delta SSE + son numéro de
// séquence monotone).
type LogEvent struct {
Seq int `json:"seq"`
Delta map[string]any `json:"delta"`
}
// Conversation est le fil unique partagé. Protégé par mu ; cond réveille les
// abonnés (aucun canal par abonné : les abonnés lisent Log au-delà de leur
// dernier Seq puis attendent cond — replay et direct sont le même chemin).
type Conversation struct {
mu sync.Mutex
cond *sync.Cond
Messages []Message `json:"messages"` // vue « modèle » (nourrit runChat)
Log []LogEvent `json:"log"` // vue « UI » rejouable
Seq int `json:"seq"`
CtxUsed int `json:"ctx_used"` // taille réelle du contexte au dernier tour
Generating bool `json:"-"`
cancel context.CancelFunc // annule la génération en cours (/stop)
epoch int // incrémenté à chaque reset → invalide les abonnés
}
var conv = func() *Conversation {
c := &Conversation{}
c.cond = sync.NewCond(&c.mu)
return c
}()
// convPath = fichier de persistance sur la box jean (JeanHome = /etc/jean en
// prod). En clair : la box déchiffre déjà pour lancer le modèle, le relais reste
// aveugle — le persister ici ne change rien à la posture E2E.
func convPath() string { return filepath.Join(JeanHome(), "conversation.json") }
// LoadConversation recharge l'état persisté au démarrage du process. Sans fichier
// (première fois) on part d'une conversation vide.
func LoadConversation() {
b, err := os.ReadFile(convPath())
if err != nil {
return
}
conv.mu.Lock()
defer conv.mu.Unlock()
_ = json.Unmarshal(b, conv)
// Une génération n'a pas pu survivre à l'arrêt du process : on repart propre.
conv.Generating = false
conv.cancel = nil
}
// persist écrit l'état sur disque (appelé en fin de tour et sur reset, pas à
// chaque delta). L'appelant NE doit PAS détenir mu.
func (c *Conversation) persist() {
c.mu.Lock()
b, err := json.Marshal(c)
c.mu.Unlock()
if err != nil {
return
}
_ = os.MkdirAll(JeanHome(), 0o755)
tmp := convPath() + ".tmp"
if err := os.WriteFile(tmp, b, 0o600); err != nil {
return
}
_ = os.Rename(tmp, convPath())
}
// appendDelta journalise un événement d'affichage et réveille les abonnés.
func (c *Conversation) appendDelta(delta map[string]any) {
c.mu.Lock()
c.Seq++
c.Log = append(c.Log, LogEvent{Seq: c.Seq, Delta: delta})
if len(c.Log) > maxLogEvents {
c.Log = c.Log[len(c.Log)-maxLogEvents:]
}
c.cond.Broadcast()
c.mu.Unlock()
}
// convState renvoie un instantané léger (pour /api/chat/state).
func (c *Conversation) state() map[string]any {
c.mu.Lock()
defer c.mu.Unlock()
return map[string]any{"seq": c.Seq, "generating": c.Generating, "ctx_used": c.CtxUsed}
}
// ErrBusy : une génération est déjà en cours (un seul tour à la fois).
var ErrBusy = fmt.Errorf("génération en cours")
// StartTurn ajoute le message utilisateur et lance la génération EN ARRIÈRE-PLAN
// (context.Background, détaché de toute connexion HTTP). Renvoie ErrBusy si un
// tour est déjà en cours, ou une erreur si le modèle n'est pas prêt.
func (c *Conversation) StartTurn(text string, caps Caps, temperature float64) error {
if !healthCheck() {
return fmt.Errorf("⏳ Le modèle est encore en train de charger — réessaie dans quelques secondes.")
}
c.mu.Lock()
if c.Generating {
c.mu.Unlock()
return ErrBusy
}
c.Generating = true
ctx, cancel := context.WithCancel(context.Background())
c.cancel = cancel
c.Messages = append(c.Messages, Message{Role: "user", Content: text})
c.mu.Unlock()
// Borne de tour + bulle utilisateur (rejouables).
c.appendDelta(map[string]any{"user": text})
if temperature == 0 {
temperature = 0.7
}
go c.generate(ctx, caps, temperature)
return nil
}
// generate exécute un tour complet et journalise chaque événement. Détaché : la
// fermeture du navigateur n'a aucun effet ici, seul /stop (cancel) l'interrompt.
func (c *Conversation) generate(ctx context.Context, caps Caps, temperature float64) {
defer func() {
c.mu.Lock()
c.Generating = false
c.cancel = nil
c.mu.Unlock()
c.appendDelta(map[string]any{"turn_done": true})
c.persist()
}()
// Snapshot de la vue modèle.
c.mu.Lock()
msgs := append([]Message(nil), c.Messages...)
ctxUsed := c.CtxUsed
c.mu.Unlock()
// Compaction proactive (façon Hermes) sur la vue MODÈLE uniquement ; le journal
// d'affichage garde le fil complet. On signale juste un toast au client.
if compacted, changed := MaybeCompact(ctx, msgs, caps, ctxUsed); changed {
msgs = compacted
c.mu.Lock()
c.Messages = compacted
c.mu.Unlock()
c.appendDelta(map[string]any{"compacted": true})
}
var content strings.Builder
extra, _ := runChat(ctx, InjectSkills(msgs, caps), temperature, caps, func(ev StreamEvent) bool {
switch {
case ev.Err != nil:
c.appendDelta(map[string]any{"error": ev.Err.Error()})
case ev.ToolUsed != nil:
c.appendDelta(map[string]any{"tool_used": map[string]any{
"name": ev.ToolUsed.Name, "label": ev.ToolUsed.Label,
"result": ev.ToolUsed.Result, "done": ev.ToolUsed.Done, "typing": ev.ToolUsed.Typing,
}})
case ev.Stats != nil:
// Taille réelle du contexte (usage.prompt_tokens + généré) pour le compteur
// et la décision de compactage au tour suivant.
if ev.Stats.PromptTokensTotal > 0 {
c.mu.Lock()
c.CtxUsed = ev.Stats.PromptTokensTotal + ev.Stats.GenTokens
c.mu.Unlock()
}
c.appendDelta(map[string]any{"stats": ev.Stats})
case ev.DropReasoning:
c.appendDelta(map[string]any{"drop_reasoning": true})
case ev.Reasoning != "":
c.appendDelta(map[string]any{"reasoning_content": ev.Reasoning})
case ev.Content != "":
content.WriteString(ev.Content)
c.appendDelta(map[string]any{"content": ev.Content})
}
return true // génération détachée : on ne s'interrompt jamais sur un abonné
})
// Persiste la vue modèle : messages d'outils (assistant tool_calls + résultats)
// PUIS la réponse finale — même ordre que l'ancien client, pour que le modèle
// garde la trace de ce qu'il a fait.
c.mu.Lock()
c.Messages = append(c.Messages, extra...)
if s := content.String(); strings.TrimSpace(s) != "" {
c.Messages = append(c.Messages, Message{Role: "assistant", Content: s})
}
c.mu.Unlock()
}
// Stop interrompt la génération en cours (le cas échéant).
func (c *Conversation) Stop() {
c.mu.Lock()
cancel := c.cancel
c.mu.Unlock()
if cancel != nil {
cancel()
}
}
// Reset démarre une nouvelle conversation (vide) pour TOUS les appareils. On
// interrompt une éventuelle génération, on vide tout et on bump epoch pour que
// les abonnés reçoivent l'ordre de nettoyer leur affichage.
func (c *Conversation) Reset() {
c.Stop()
c.mu.Lock()
c.Messages = nil
c.Log = nil
c.Seq = 0
c.CtxUsed = 0
c.epoch++
c.cond.Broadcast()
c.mu.Unlock()
c.persist()
}
// Subscribe diffuse les événements au client via emit, en commençant par le
// replay de Log[from:] puis en suivant le direct. Bloque jusqu'à ce que ctx (la
// connexion HTTP) soit annulé — la génération, elle, continue indépendamment.
// emit renvoie false si l'écriture échoue (client parti) → on sort.
func (c *Conversation) Subscribe(ctx context.Context, from int, emit func(map[string]any) bool) {
// Réveille les attentes de cond quand la connexion se ferme (sinon cond.Wait
// resterait bloqué faute de nouvel événement).
go func() {
<-ctx.Done()
c.mu.Lock()
c.cond.Broadcast()
c.mu.Unlock()
}()
c.mu.Lock()
last := from
epoch := c.epoch
for {
if ctx.Err() != nil {
c.mu.Unlock()
return
}
// Un reset a eu lieu : on ordonne au client de nettoyer et on repart de 0.
if c.epoch != epoch {
epoch = c.epoch
last = 0
c.mu.Unlock()
if !emit(map[string]any{"reset": true}) {
return
}
c.mu.Lock()
continue
}
// Envoie tout ce qui est plus récent que `last`.
sent := false
for _, ev := range c.Log {
if ev.Seq <= last {
continue
}
last = ev.Seq
delta := ev.Delta
c.mu.Unlock()
out := map[string]any{"seq": ev.Seq}
for k, v := range delta {
out[k] = v
}
if !emit(out) {
return
}
c.mu.Lock()
sent = true
if ctx.Err() != nil {
c.mu.Unlock()
return
}
}
if sent {
continue // il peut y avoir eu de nouveaux événements pendant l'envoi
}
c.cond.Wait() // rien de neuf : dort jusqu'au prochain Broadcast (ou ctx annulé)
}
}
+124
View File
@@ -0,0 +1,124 @@
package main
import (
"context"
"net/http"
"net/http/httptest"
"strings"
"sync"
"testing"
"time"
)
// newTestConv crée une conversation isolée (le global `conv` sert au process réel).
func newTestConv() *Conversation {
c := &Conversation{}
c.cond = sync.NewCond(&c.mu)
return c
}
// Subscribe doit rejouer les événements déjà journalisés PUIS suivre le direct,
// et se terminer quand le contexte (la connexion) est annulé.
func TestSubscribeReplayAndLive(t *testing.T) {
c := newTestConv()
c.appendDelta(map[string]any{"user": "salut"})
c.appendDelta(map[string]any{"content": "bon"})
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
got := make(chan int, 32)
go c.Subscribe(ctx, 0, func(m map[string]any) bool {
if s, ok := m["seq"].(int); ok {
got <- s
}
return true
})
// Les 2 événements déjà présents (replay).
waitSeq(t, got, 1)
waitSeq(t, got, 2)
// Un événement en direct après abonnement.
c.appendDelta(map[string]any{"content": "jour"})
waitSeq(t, got, 3)
}
// Un abonné qui démarre à from=N ne reçoit que ce qui est plus récent que N.
func TestSubscribeFromOffset(t *testing.T) {
c := newTestConv()
c.appendDelta(map[string]any{"user": "a"})
c.appendDelta(map[string]any{"user": "b"})
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
got := make(chan int, 8)
go c.Subscribe(ctx, 1, func(m map[string]any) bool {
if s, ok := m["seq"].(int); ok {
got <- s
}
return true
})
// from=1 → on saute le seq 1, on reçoit 2 en premier.
waitSeq(t, got, 2)
}
// Reset bump l'epoch et pousse un {reset:true} aux abonnés, qui repartent de 0.
func TestResetNotifiesSubscribers(t *testing.T) {
c := newTestConv()
c.appendDelta(map[string]any{"user": "x"})
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
resetSeen := make(chan bool, 4)
go c.Subscribe(ctx, 0, func(m map[string]any) bool {
if _, ok := m["reset"]; ok {
resetSeen <- true
}
return true
})
// laisse l'abonné consommer le replay initial
time.Sleep(20 * time.Millisecond)
c.Reset()
select {
case <-resetSeen:
case <-time.After(time.Second):
t.Fatal("l'abonné n'a pas reçu l'événement reset")
}
if c.Seq != 0 || len(c.Log) != 0 {
t.Fatalf("après Reset: Seq=%d len(Log)=%d, attendu 0/0", c.Seq, len(c.Log))
}
}
// Les handlers de contrôle répondent en JSON sans dépendre du modèle.
func TestChatControlHandlers(t *testing.T) {
// reset → conversation vide
rr := httptest.NewRecorder()
handleChatReset(rr, httptest.NewRequest("POST", "/api/chat/reset", nil))
if rr.Code != 200 {
t.Fatalf("reset code %d", rr.Code)
}
// state → seq 0, pas de génération
rr = httptest.NewRecorder()
handleChatState(rr, httptest.NewRequest("GET", "/api/chat/state", nil))
if rr.Code != 200 || !strings.Contains(rr.Body.String(), "\"seq\":0") {
t.Fatalf("state inattendu: %d %s", rr.Code, rr.Body.String())
}
// send message vide → 400
rr = httptest.NewRecorder()
req := httptest.NewRequest("POST", "/api/chat/send", strings.NewReader(`{"message":""}`))
req.Header.Set("Content-Type", "application/json")
handleChatSend(rr, req)
if rr.Code != http.StatusBadRequest {
t.Fatalf("send vide devrait être 400, obtenu %d", rr.Code)
}
}
func waitSeq(t *testing.T, ch <-chan int, want int) {
t.Helper()
select {
case s := <-ch:
if s != want {
t.Fatalf("seq reçu %d, attendu %d", s, want)
}
case <-time.After(time.Second):
t.Fatalf("timeout en attendant seq %d", want)
}
}
+104 -207
View File
@@ -1317,11 +1317,16 @@ function addCopyButtons(root){
pre.appendChild(btn);
});
}
function resetChat(){ msgs=[]; document.getElementById('chat').innerHTML=''; saveChat(); try{localStorage.removeItem('jean.ctxused');}catch(e){} setCtxUsed(0); toast('chat vidé'); }
// Nouvelle conversation POUR TOUS LES APPAREILS : le serveur vide le fil et
// diffuse un {reset} ; le flux d'abonnement nettoie alors l'affichage.
function resetChat(){ jfetch('/api/chat/reset',{method:'POST'}).catch(()=>{}); toast('nouvelle conversation'); }
// Compaction : on demande à l'IA un résumé de la conversation destiné à la
// reprendre dans une session neuve, puis on repart d'un contexte propre seedé
// avec ce résumé. Réduit drastiquement les tokens tout en gardant le fil.
async function compactContext(){
// Le compactage est désormais AUTOMATIQUE côté serveur (façon Hermes) : ce
// bouton manuel n'a plus d'objet et est neutralisé.
toast('le compactage est automatique (façon Hermes)'); return;
if(busy) return;
if(!msgs.length){ toast('rien à compacter'); return; }
if(!await askConfirm('L\'IA va résumer la conversation puis repartir d\'une session propre basée sur ce résumé.', {title:'Compacter le contexte ?', okText:'Compacter'})) return;
@@ -1364,218 +1369,110 @@ async function compactContext(){
// Persistance de la conversation : on garde user+assistant en localStorage pour
// survivre à un refresh (les bulles tool/reasoning sont éphémères, non stockées).
function saveChat(){ try{ localStorage.setItem('jean.chat', JSON.stringify(msgs)); }catch(e){} }
function restoreChat(){
let raw=null; try{ raw=localStorage.getItem('jean.chat'); }catch(e){}
if(!raw) return;
try{ msgs=JSON.parse(raw)||[]; }catch(e){ msgs=[]; return; }
for(const m of msgs){
// Les messages d'outil (role 'tool' et assistant porteur de tool_calls sans
// texte) font partie de l'historique envoyé au modèle mais ne s'affichent pas :
// on saute pour ne pas créer de bulles vides au rechargement.
if(m.role==='tool') continue;
if(m.role==='user') addMsg('user', m.content);
else if(m.role==='assistant'){
if(!m.content && m.tool_calls) continue;
const el=addMsg('assistant',''); renderBody(el, m.content||'');
}
}
const cu=parseInt(localStorage.getItem('jean.ctxused')||'0',10); if(cu>0) setCtxUsed(cu);
}
// Source de vérité = SERVEUR. Au chargement on ouvre le flux d'abonnement
// permanent (connectStream), qui rejoue tout le fil depuis le serveur — texte,
// appels d'outils, vitesses, raisonnement — puis suit le direct. Plus de
// localStorage : le même contexte est partagé par tous les appareils.
function restoreChat(){ connectStream(); }
function onKey(e){ if(e.key==='Enter' && !e.shiftKey){ e.preventDefault(); send(); } }
let abortCtrl=null;
function stopGen(){ if(abortCtrl){ abortCtrl.abort(); toast('stop'); } }
// ===== Conversation SERVEUR (source de vérité, partagée entre appareils) =====
// L'historique et la génération vivent sur le serveur jean. Le client ouvre un
// flux d'ABONNEMENT permanent (SSE) qui rejoue le journal depuis lastSeq puis
// suit le direct. Fermer l'onglet n'arrête plus la génération (détachée côté
// serveur) ; se reconnecter rejoue tout le fil, détails compris.
let lastSeq=0, streamAbort=null;
// État de rendu du tour courant, délimité par les événements user / turn_done.
let T=null;
function newTurn(){ T={ reasonEl:null, contentEl:null, pendingToolEl:null, typingEl:null, fullContent:'', fullReason:'', turnCollapsibles:[], serverStats:null }; }
newTurn();
const simpleMode=()=>document.documentElement.getAttribute('data-display')==='simple';
function removeTyping(){ if(T.typingEl){ T.typingEl.remove(); T.typingEl=null; } }
function killTyping(){ if(!T.typingEl) return; if(simpleMode()) return; removeTyping(); }
function showTyping(){ if(!simpleMode()) return; const c=document.getElementById('chat');
if(!T.typingEl){ T.typingEl=addTyping(); } else if(c.lastElementChild!==T.typingEl){ c.appendChild(T.typingEl); } }
// Label de vitesse rendu depuis les valeurs (fonctionne aussi bien en direct
// qu'au replay — pas de timer performance.now, qui n'a pas de sens hors-ligne).
function renderStats(el, s){
if(!el||!s) return;
const role=(el===T.contentEl)?'assistant':'reasoning';
const parts=[role];
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');
if(s.gen_tokens) parts.push('decode '+s.gen_tokens+' tok · '+(s.gen_per_second||0).toFixed(1)+' tok/s');
setLabel(el, parts.join(' · '));
}
function setBusy(on){ busy=on;
document.getElementById('send').style.display=on?'none':'inline-block';
document.getElementById('stop').style.display=on?'inline-block':'none'; }
// Traite UN événement du flux — même sémantique que l'ancien switch inline, mais
// piloté par le serveur et rejouable à l'identique.
function handleDelta(d){
if(typeof d.seq==='number' && d.seq>lastSeq) lastSeq=d.seq;
if(d.reset!==undefined){ document.getElementById('chat').innerHTML=''; newTurn(); setCtxUsed(0); lastSeq=0; setBusy(false); return; }
if(d.user!==undefined){ newTurn(); addMsg('user', d.user); setBusy(true); T.typingEl=addTyping(); return; }
if(d.turn_done){ removeTyping(); collapseAll(T.turnCollapsibles); if(T.serverStats) renderStats(T.contentEl||T.reasonEl, T.serverStats); setBusy(false); return; }
if(d.error){ removeTyping(); T.contentEl=null; T.reasonEl=null; const eb=addMsg('assistant',''); eb.classList.add('errmsg'); renderBody(eb, d.error); return; }
if(d.compacted){ toast('contexte compacté — les vieux tours ont été résumés'); return; }
if(d.stats){ T.serverStats=d.stats;
if(d.stats.prompt_tokens_total){ setCtxUsed((d.stats.prompt_tokens_total||0)+(d.stats.gen_tokens||0)); }
if(T.contentEl||T.reasonEl) renderStats(T.contentEl||T.reasonEl, d.stats); return; }
if(d.tool_used){
killTyping(); T.contentEl=null; T.reasonEl=null; const tu=d.tool_used;
if(!T.pendingToolEl){ collapseAll(T.turnCollapsibles); T.pendingToolEl=addMsg('tool',''); T.turnCollapsibles.push(T.pendingToolEl); }
renderToolMsg(T.pendingToolEl, tu);
if(!tu.done) showTyping();
if(tu.done){ T.pendingToolEl=null; if(tu.name==='mem_add'||tu.name==='mem_edit') loadMem(); }
return; }
if(d.drop_reasoning){
if(T.reasonEl){ const i=T.turnCollapsibles.indexOf(T.reasonEl); if(i>=0) T.turnCollapsibles.splice(i,1); T.reasonEl.remove(); T.reasonEl=null; T.fullReason=''; }
return; }
if(d.reasoning_content){
killTyping();
if(!T.reasonEl){ collapseAll(T.turnCollapsibles); T.reasonEl=addMsg('reasoning',''); T.fullReason=''; T.turnCollapsibles.push(T.reasonEl); }
showTyping(); T.fullReason+=d.reasoning_content; renderBody(T.reasonEl, T.fullReason);
return; }
if(d.content){
removeTyping();
if(!T.contentEl){ collapseAll(T.turnCollapsibles); T.contentEl=addMsg('assistant',''); T.fullContent=''; }
T.fullContent+=d.content; renderBody(T.contentEl, T.fullContent);
return; }
}
// Flux d'abonnement permanent + reconnexion auto (from=lastSeq → pas de
// re-téléchargement complet après une coupure / bascule d'appareil).
async function connectStream(){
while(true){
streamAbort=new AbortController();
try{
const r=await jfetch('/api/chat',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({from:lastSeq}),signal:streamAbort.signal});
const reader=r.body.getReader(); const dec=new TextDecoder(); let buf='';
while(true){
const {done,value}=await reader.read(); if(done) break;
buf+=dec.decode(value,{stream:true}); let i;
while((i=buf.indexOf('\n\n'))>=0){
const chunk=buf.slice(0,i); buf=buf.slice(i+2);
for(const line of chunk.split('\n')){
if(!line.startsWith('data:')) continue;
const data=line.slice(5).trim(); if(data===''||data==='[DONE]') continue;
try{ const o=JSON.parse(data); const d=(o.choices&&o.choices[0]&&o.choices[0].delta)||{}; handleDelta(d); }catch(e){}
}
}
}
}catch(e){ /* coupure : on reconnecte */ }
await new Promise(res=>setTimeout(res, 1200));
}
}
function stopGen(){ jfetch('/api/chat/stop',{method:'POST'}).catch(()=>{}); toast('stop'); }
async function send(){
if(busy) return;
const ta=document.getElementById('input'); const text=ta.value.trim();
if(!text) return;
ta.value=''; busy=true;
document.getElementById('send').style.display='none';
document.getElementById('stop').style.display='inline-block';
addMsg('user', text);
msgs.push({role:'user',content:text}); saveChat();
let typingEl=addTyping();
const simpleMode=()=>document.documentElement.getAttribute('data-display')==='simple';
// Retire définitivement l'indicateur de frappe (fin de tour).
const removeTyping=()=>{ if(typingEl){ typingEl.remove(); typingEl=null; } };
// En mode COMPLET : on retire l'indicateur dès qu'une bulle prend le relais.
// En mode SIMPLE : raisonnement/outils sont masqués, donc on GARDE le même
// indicateur (mêmes points animés) vivant tout le tour, remis en bas.
// COMPLET : retire l'indicateur dès qu'une bulle prend le relais.
// SIMPLE : on le GARDE (raisonnement/outils masqués) — sans le toucher, pour
// ne PAS redémarrer son animation CSS (déplacer le nœud la remettrait à zéro).
const killTyping=()=>{ if(!typingEl) return; if(simpleMode()) return; removeTyping(); };
// SIMPLE : garantit l'indicateur (le recrée s'il manque) et le remet en bas
// UNIQUEMENT s'il n'y est pas déjà — idempotent, donc l'animation ne se
// réinitialise pas à chaque token pendant que l'IA raisonne / appelle un outil.
const showTyping=()=>{ if(!simpleMode()) return; const c=document.getElementById('chat');
if(!typingEl){ typingEl=addTyping(); } else if(c.lastElementChild!==typingEl){ c.appendChild(typingEl); } };
let reasonEl=null, contentEl=null, pendingToolEl=null, fullContent='', fullReason='', toolMsgs=null, compacted=false;
let turnCollapsibles=[]; // bulles reasoning/tool à replier dès la réponse finale
let tokCount=0, reasonTokCount=0, t0=performance.now(), firstTok=0, statsTimer=null;
let serverStats=null, streamErr=false;
const updateStats=(final)=>{
if(!contentEl && !reasonEl) return;
const el = contentEl || reasonEl;
const role = contentEl ? 'assistant' : 'reasoning';
// Mid-stream: client-side estimate. Final tick: use the server's exact
// timings if it gave us any (more accurate, includes prefill).
if(final && serverStats){
const s = serverStats;
const totalSec = (performance.now() - t0) / 1000;
const parts = [role];
// prefill = tokens RÉELLEMENT traités ce tour (timings.prompt_n) : gros au 1er
// tour (system prompt), petit ensuite car le préfixe est en cache. On NE prend
// PAS prompt_tokens_total (taille totale du prompt ≈ tout le contexte à chaque
// tour) — ça donnait un prefill quasi constant, incohérent avec la vitesse
// affichée (prompt_per_second, qui mesure justement prompt_n). Repli sur le
// total seulement si le backend n'a pas envoyé de timings.
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');
if(s.gen_tokens) parts.push('decode ' + s.gen_tokens + ' tok · ' + s.gen_per_second.toFixed(1) + ' tok/s');
parts.push('total ' + totalSec.toFixed(1) + 's');
setLabel(el, parts.join(' · '));
return;
}
const n = contentEl ? tokCount : reasonTokCount;
const elapsed = (performance.now() - (firstTok||t0)) / 1000;
const tps = elapsed > 0.05 ? (n / elapsed) : 0;
setLabel(el, role + ' · ' + n + ' tok · ' + tps.toFixed(1) + ' tok/s');
};
abortCtrl = new AbortController();
ta.value='';
try{
const sys=(document.getElementById('sysprompt').value||'').trim();
const out = sys ? [{role:'system',content:sys}, ...msgs] : msgs;
const r=await jfetch('/api/chat',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({messages:out,ctx_used:CTX_USED}),signal:abortCtrl.signal});
const reader=r.body.getReader(); const dec=new TextDecoder(); let buf='';
statsTimer = setInterval(updateStats, 250);
while(true){
const {done,value}=await reader.read(); if(done) break;
buf+=dec.decode(value,{stream:true});
let i;
while((i=buf.indexOf('\n\n'))>=0){
const chunk=buf.slice(0,i); buf=buf.slice(i+2);
for(const line of chunk.split('\n')){
if(!line.startsWith('data:')) continue;
const data=line.slice(5).trim();
if(data==='[DONE]') continue;
try{
const o=JSON.parse(data);
const d=(o.choices&&o.choices[0]&&o.choices[0].delta)||{};
if(d.error){
removeTyping();
// Erreur serveur (ex : modèle en cours de chargement → 503). On
// l'affiche, on retire le message user de l'historique et on le
// remet dans l'input pour renvoyer d'un clic une fois prêt.
contentEl=null; reasonEl=null;
const eb=addMsg('assistant',''); eb.classList.add('errmsg');
renderBody(eb, d.error);
if(msgs.length && msgs[msgs.length-1].role==='user'){ ta.value=msgs[msgs.length-1].content; msgs.pop(); }
saveChat();
streamErr=true;
try{ reader.cancel(); }catch(_){}
break;
}
if(d.stats){ serverStats = d.stats; }
// Tool-turn messages (assistant tool_calls + tool results) sent at the
// end of the stream. On les garde pour les réinjecter dans msgs AVANT
// la réponse finale, sinon le modèle re-déclenche le même skill/commande
// à chaque tour (aucune trace de l'avoir déjà fait).
if(d.tool_messages){ toolMsgs = d.tool_messages; }
if(d.history_replace){
// Compaction côté serveur (façon Hermes) : on remplace notre
// historique par la version résumée. msgs ne contient jamais le
// system prompt (ajouté seulement à l'envoi), donc on retire les
// éventuels messages système de tête pour ne pas les dupliquer.
let h=d.history_replace.slice();
while(h.length && h[0].role==='system') h.shift();
msgs=h; saveChat(); compacted=true;
}
if(d.tool_used){
killTyping();
// Close out the current assistant/reasoning bubbles so any text the
// model streams *after* the tool call starts a fresh bubble BELOW it
// (and isn't appended above, into the pre-call bubble).
contentEl=null; reasonEl=null;
const tu=d.tool_used;
// One bubble per call: the live "typing" updates, the "en cours"
// spinner, and the final result all re-render the SAME element.
// Chaînage : replie les étapes précédentes (reasoning/outils déjà
// finis) dès qu'une nouvelle démarre — seule l'étape en cours reste ouverte.
if(!pendingToolEl){ collapseAll(turnCollapsibles); pendingToolEl=addMsg('tool',''); turnCollapsibles.push(pendingToolEl); }
renderToolMsg(pendingToolEl, tu);
if(!tu.done) showTyping(); // simple : points animés pendant l'appel d'outil
if(tu.done){
pendingToolEl=null;
// L'IA gère sa mémoire elle-même : rafraîchir le menu en direct
// dès qu'une page est créée/modifiée (pas besoin de F5).
if(tu.name==='mem_add' || tu.name==='mem_edit') loadMem();
}
}
if(d.drop_reasoning){
// Le modèle a « pensé sans agir » : on relance le tour. Le
// raisonnement mort-né est retiré pour ne pas afficher un double bloc.
if(reasonEl){
const i=turnCollapsibles.indexOf(reasonEl); if(i>=0) turnCollapsibles.splice(i,1);
reasonEl.remove(); reasonEl=null; fullReason='';
}
}
if(d.reasoning_content){
killTyping();
// Start a fresh buffer per bubble — otherwise a reasoning block that
// follows a tool call re-renders ALL previously accumulated reasoning
// into the new bubble (duplicated dozens of times over a turn).
if(!reasonEl){ collapseAll(turnCollapsibles); reasonEl=addMsg('reasoning',''); fullReason=''; turnCollapsibles.push(reasonEl); }
showTyping(); // simple : points animés pendant que l'IA raisonne
if(!firstTok) firstTok=performance.now();
reasonTokCount++;
fullReason+=d.reasoning_content; renderBody(reasonEl, fullReason);
}
if(d.content){
removeTyping(); // la réponse prend le relais : plus de points (comme en complet)
if(!contentEl){ collapseAll(turnCollapsibles); contentEl=addMsg('assistant',''); fullContent=''; firstTok=performance.now(); }
tokCount++;
fullContent+=d.content; renderBody(contentEl, fullContent);
}
}catch(e){}
}
if(streamErr) break;
}
if(streamErr) break;
}
if(!streamErr){
if(toolMsgs) msgs.push(...toolMsgs);
msgs.push({role:'assistant',content:fullContent}); saveChat();
}
if(compacted) toast('contexte compacté — les vieux tours ont été résumés');
// Le KV cache s'accumule : chaque tour ajoute les tokens NOUVELLEMENT traités
// (prompt_tokens, le préfixe caché n'est pas recompté) + ceux générés. La
// taille réelle du contexte = somme cumulée, donc monotone croissante.
if(serverStats){
// prompt_tokens_total (usage) = taille TOTALE du prompt ce tour (system +
// historique, préfixe caché compris) → contexte ABSOLU et exact. Sinon repli
// sur l'ancienne accumulation (prompt_tokens = seulement le neuf traité).
if(serverStats.prompt_tokens_total){ setCtxUsed((serverStats.prompt_tokens_total||0) + (serverStats.gen_tokens||0)); }
else { setCtxUsed(CTX_USED + (serverStats.prompt_tokens||0) + (serverStats.gen_tokens||0)); }
try{ localStorage.setItem('jean.ctxused', CTX_USED); }catch(e){}
}
}catch(e){
if(e.name==='AbortError'){
if(contentEl) setLabel(contentEl, 'assistant · '+tokCount+' tok · interrompu');
if(toolMsgs) msgs.push(...toolMsgs);
if(fullContent) msgs.push({role:'assistant',content:fullContent});
else if(!toolMsgs) msgs.pop();
} else {
addMsg('assistant','[erreur] '+e.message); msgs.pop();
}
saveChat();
}
removeTyping();
if(statsTimer) clearInterval(statsTimer);
updateStats(true);
abortCtrl=null; busy=false;
document.getElementById('send').style.display='inline-block';
document.getElementById('stop').style.display='none';
document.getElementById('input').focus();
const r=await jfetch('/api/chat/send',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({message:text,ctx_used:CTX_USED})});
if(r.status===409){ toast('génération déjà en cours'); ta.value=text; return; }
if(!r.ok){ let m='erreur'; try{ m=(await r.json()).error||m; }catch(_){} toast(m); ta.value=text; return; }
// La bulle utilisateur et les tokens arrivent par le flux d'abonnement.
}catch(e){ toast('échec envoi: '+e.message); ta.value=text; }
}
loadAll();
setInterval(loadStatus, 5000);
+73 -58
View File
@@ -57,7 +57,12 @@ func cmdWeb(args []string) error {
// newWebMux construit le routeur HTTP de l'UI web. Extrait de cmdWeb pour être
// réutilisé par `jean link`, qui sert ce même mux à travers le tunnel sans
// repasser par un écouteur TCP local.
var convLoadOnce sync.Once
func newWebMux() *http.ServeMux {
// Charge l'état de conversation persisté (une fois par process : jean web ET
// jean link serve appellent newWebMux).
convLoadOnce.Do(LoadConversation)
mux := http.NewServeMux()
// Pages publiques : le HTML et le JS ne contiennent aucun secret. Toute la
// donnée et toutes les actions passent par /api/* qui, lui, exige la clé.
@@ -111,8 +116,12 @@ func newWebMux() *http.ServeMux {
api("/api/restart", svcHandler("restart"))
api("/api/bench", handleBench)
api("/api/bench/last", handleBenchLast)
api("/api/chat", handleChat)
api("/api/e2e/chat", handleE2EChat) // chat chiffré E2E (boîte noire via le relais)
api("/api/chat", handleChat) // flux d'ABONNEMENT (SSE) : rejoue + suit le fil
api("/api/chat/send", handleChatSend) // envoie un message (lance la génération détachée)
api("/api/chat/stop", handleChatStop) // interrompt la génération en cours
api("/api/chat/reset", handleChatReset) // nouvelle conversation (pour tous les appareils)
api("/api/chat/state", handleChatState) // instantané léger {seq, generating, ctx_used}
api("/api/e2e/chat", handleE2EChat) // même flux mais chiffré E2E (boîte noire via le relais)
return mux
}
@@ -787,6 +796,27 @@ type chatReq struct {
// générés), rapportée par le client qui l'affiche déjà. Sert à décider du
// compactage sur le VRAI décompte plutôt qu'une estimation. 0 = inconnu.
CtxUsed int `json:"ctx_used"`
// Nouveau modèle « conversation serveur » : Message = texte du tour à lancer
// (via /api/chat/send) ; From = dernier Seq déjà vu par le client (le flux
// d'abonnement rejoue Log[From:] puis suit le direct).
Message string `json:"message"`
From int `json:"from"`
}
// capsFromBody dérive les capacités du tour à partir des overrides éventuels du
// corps de requête (agents ajean.link portant leurs propres toggles), sinon la
// config machine.
func capsFromBody(body chatReq) Caps {
caps := globalCaps()
if body.Agent != nil {
caps.Agent = *body.Agent
} else if body.Tools != nil || body.Skills != nil {
caps.Agent = (body.Tools != nil && *body.Tools) || (body.Skills != nil && *body.Skills)
}
if body.Internet != nil {
caps.Internet = *body.Internet && crawlReachable()
}
return caps
}
// sseHeartbeat garde la réponse SSE active en écrivant un commentaire (`: ping`,
@@ -822,67 +852,52 @@ func sseHeartbeat(w http.ResponseWriter, flusher http.Flusher) (*sync.Mutex, fun
return mu, func() { close(done) }
}
// runChatStream exécute le chat et pousse chaque événement (delta) via emit, qui
// renvoie false pour interrompre. Partagé par handleChat et handleE2EChat.
// 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) {
// Garde-fou : si le modèle n'est pas encore chargé, llama-server répond 503
// ("loading model") et le tour partirait dans le vide (aucune réponse, l'user
// renvoie en boucle). On renvoie une erreur explicite affichée dans le chat.
if !healthCheck() {
emit(map[string]any{"error": "⏳ Le modèle est encore en train de charger — patiente quelques secondes puis renvoie ton message."})
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
}
if body.Temperature == 0 {
body.Temperature = 0.7
if strings.TrimSpace(body.Message) == "" {
sendJSON(w, 400, map[string]any{"ok": false, "error": "message vide"})
return
}
caps := globalCaps()
if body.Agent != nil {
caps.Agent = *body.Agent
} else if body.Tools != nil || body.Skills != nil {
// rétro-compat : anciens clients qui envoyaient deux drapeaux séparés
caps.Agent = (body.Tools != nil && *body.Tools) || (body.Skills != nil && *body.Skills)
}
if body.Internet != nil {
// On garde la cohérence prompt/outils : internet demandé ET serveur joignable.
caps.Internet = *body.Internet && crawlReachable()
}
// Compaction proactive (façon Hermes) : si l'historique dépasse le seuil, on
// le résume AVANT le tour et on renvoie l'historique compacté au client pour
// qu'il remplace le sien — évite de re-renvoyer tout à chaque tour.
hist := body.Messages
if compacted, changed := MaybeCompact(ctx, hist, caps, body.CtxUsed); changed {
hist = compacted
emit(map[string]any{"history_replace": hist})
}
msgs := InjectSkills(hist, caps)
extra, _ := runChat(ctx, msgs, body.Temperature, caps, func(ev StreamEvent) bool {
if ev.Err != nil {
return emit(map[string]any{"error": ev.Err.Error()})
}
if ev.ToolUsed != nil {
return emit(map[string]any{"tool_used": map[string]any{"name": ev.ToolUsed.Name, "label": ev.ToolUsed.Label, "result": ev.ToolUsed.Result, "done": ev.ToolUsed.Done, "typing": ev.ToolUsed.Typing}})
}
if ev.Stats != nil {
return emit(map[string]any{"stats": ev.Stats})
}
if ev.DropReasoning {
return emit(map[string]any{"drop_reasoning": true})
}
if ev.Reasoning != "" {
return emit(map[string]any{"reasoning_content": ev.Reasoning})
}
if ev.Content != "" {
return emit(map[string]any{"content": ev.Content})
}
return true
})
// Surface the tool-turn messages (assistant tool_calls + tool results) as a
// final event so the stateless web client can store them in its history,
// BEFORE the final assistant text. Without this the browser only keeps the
// final answer and the model re-invokes the same skill/command every turn.
if len(extra) > 0 {
emit(map[string]any{"tool_messages": extra})
if err := conv.StartTurn(body.Message, 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})
}
func handleChatState(w http.ResponseWriter, r *http.Request) {
sendJSON(w, 200, conv.state())
}
func handleChat(w http.ResponseWriter, r *http.Request) {