diff --git a/conversation.go b/conversation.go new file mode 100644 index 0000000..73344eb --- /dev/null +++ b/conversation.go @@ -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é) + } +} diff --git a/conversation_test.go b/conversation_test.go new file mode 100644 index 0000000..3cbdcf0 --- /dev/null +++ b/conversation_test.go @@ -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) + } +} diff --git a/ui/index.html b/ui/index.html index bbdcb24..b71603c 100644 --- a/ui/index.html +++ b/ui/index.html @@ -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); diff --git a/web.go b/web.go index f17d3cd..9ffa2b3 100644 --- a/web.go +++ b/web.go @@ -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) {