diff --git a/internal/loki/chat_conversation.go b/internal/loki/chat_conversation.go index 90b5ee5..1333d7d 100644 --- a/internal/loki/chat_conversation.go +++ b/internal/loki/chat_conversation.go @@ -68,6 +68,27 @@ type Conversation struct { // à une génération fantôme dans la discussion ouverte. runningTaskID string runningTaskName string + + // File d'attente des messages envoyés PENDANT une génération (AJEAN + // 0.14.0). Ils sont injectés dans le tour en cours à la prochaine frontière + // d'étape (après un appel d'outil), ou — si le tour se termine avant — + // traités comme tours suivants, dans l'ordre. Protégée par mu. + queued []queuedMsg + // recentCIDs : identifiants d'envoi déjà acceptés. L'UI réessaie un envoi + // dont la réponse s'est perdue (tunnel) : avant la file, le 409 « déjà en + // cours » faisait office de dédoublonnage ; sans lui, le réessai mettrait + // le même message deux fois en file. + recentCIDs []string +} + +// queuedMsg = un message mis en file pendant la génération. Caps et +// température sont ceux choisis à l'envoi : un tour démarré depuis la file doit +// s'exécuter avec ces réglages-là. +type queuedMsg struct { + text string + files []attachInfo + caps Caps + temp float64 } var conv = func() *Conversation { @@ -167,6 +188,7 @@ func (c *Conversation) loadFrom(b []byte) { c.Stop() c.mu.Lock() c.Messages, c.Log, c.Seq, c.CtxUsed = nil, nil, 0, 0 + c.queued = nil // file de l'ancienne discussion : elle ne suit pas la bascule if len(b) > 0 { _ = json.Unmarshal(b, c) c.Messages = stripImageParts(c.Messages) // même guérison qu'au chargement @@ -506,6 +528,131 @@ func (c *Conversation) StartTurn(text string, files []attachInfo, caps Caps, tem return nil } +// ErrDupSend : cet envoi (même cid) a déjà été accepté — réessai d'un client +// qui n'a pas reçu la réponse. Ce n'est pas une erreur pour l'utilisateur. +var ErrDupSend = fmt.Errorf("envoi déjà reçu") + +// seenCID enregistre cid et dit s'il avait déjà été vu. Appelant : mu tenu. +func (c *Conversation) seenCIDLocked(cid string) bool { + if cid == "" { + return false + } + for _, v := range c.recentCIDs { + if v == cid { + return true + } + } + c.recentCIDs = append(c.recentCIDs, cid) + if len(c.recentCIDs) > 64 { + c.recentCIDs = c.recentCIDs[len(c.recentCIDs)-64:] + } + return false +} + +// EnqueueOrStart démarre un tour tout de suite si le moteur est libre, sinon +// MET EN FILE le message (AJEAN 0.14.0) — au lieu du 409 d'avant, qui obligeait +// à arrêter la réponse pour ajouter une précision. Seul un tour de CHAT accepte +// une file : une tâche planifiée qui occupe le modèle garde le refus +// (busyReason), sa fin ne dépile rien. cid = identifiant de l'envoi (voir +// recentCIDs) ; un doublon renvoie ErrDupSend. +func (c *Conversation) EnqueueOrStart(cid, text string, files []attachInfo, caps Caps, temperature float64) (bool, error) { + c.mu.Lock() + if c.seenCIDLocked(cid) { + c.mu.Unlock() + return false, ErrDupSend + } + if c.Generating && c.runningTaskName == "" { + c.queued = append(c.queued, queuedMsg{text: text, files: files, caps: caps, temp: temperature}) + c.mu.Unlock() + return true, nil + } + c.mu.Unlock() + err := c.StartTurn(text, files, caps, temperature) + if err != nil { + // Refusé (modèle pas prêt, tâche en cours) : le client pourra renvoyer + // ce même message, son cid ne doit pas rester marqué comme reçu. + c.mu.Lock() + for i, v := range c.recentCIDs { + if v == cid { + c.recentCIDs = append(c.recentCIDs[:i], c.recentCIDs[i+1:]...) + break + } + } + c.mu.Unlock() + } + return false, err +} + +// drainQueued vide la file et renvoie les messages utilisateur à injecter dans +// le tour EN COURS (vue modèle), en les journalisant pour qu'ils apparaissent +// dans le fil sur tous les appareils. Appelé par runChat entre deux étapes. +// epoch garde contre un reset survenu entre-temps. +func (c *Conversation) drainQueued(epoch int) []Message { + c.mu.Lock() + if c.epoch != epoch || len(c.queued) == 0 { + c.mu.Unlock() + return nil + } + items := c.queued + c.queued = nil + c.mu.Unlock() + + var out []Message + for _, q := range items { + delta := map[string]any{"user": q.text} + if m := chatModelName(); m != "" { + delta["model"] = m + } + if len(q.files) > 0 { + delta["files"] = q.files + } + c.appendDelta(epoch, delta) + prompt := q.text + if strings.TrimSpace(prompt) == "" { + prompt = "Prends-en connaissance." + } + out = append(out, Message{Role: "user", Content: userMessageContent(q.files, prompt)}) + } + return out +} + +// startQueuedIfAny démarre le prochain tour depuis la file, s'il en reste et +// que le moteur est libre. Appelé à la toute fin d'un tour : les messages +// envoyés après la dernière frontière d'étape sont ainsi traités sans que +// l'utilisateur ait à les renvoyer. Chaque tour lancé ainsi rappelle +// startQueuedIfAny à sa fin → la file se vide dans l'ordre. Renvoie true si un +// tour a démarré. +func (c *Conversation) startQueuedIfAny() bool { + c.mu.Lock() + if c.Generating || len(c.queued) == 0 { + c.mu.Unlock() + return false + } + q := c.queued[0] + c.queued = c.queued[1:] + epoch := c.epoch + c.mu.Unlock() + if err := c.StartTurn(q.text, q.files, q.caps, q.temp); err != nil { + // Modèle indisponible au moment de dépiler (rare) : on le dit dans le + // fil plutôt que de perdre le message en silence, et on tente le suivant. + c.appendDelta(epoch, map[string]any{"error": "message en attente non envoyé : " + err.Error()}) + return c.startQueuedIfAny() + } + return true +} + +// dropQueued abandonne la file (bouton stop : l'utilisateur reprend la main) et +// le signale au client, qui retire ses messages « en attente ». +func (c *Conversation) dropQueued(epoch int) { + c.mu.Lock() + n := len(c.queued) + c.queued = nil + c.mu.Unlock() + if n > 0 { + c.appendDelta(epoch, map[string]any{"queue_dropped": n}) + } +} + // chatModelName renvoie le nom lisible du modèle configuré : fichier sans // dossier ni extension .gguf (même règle que modelLabel côté UI), "" si aucun. func chatModelName() string { @@ -550,15 +697,24 @@ func (c *Conversation) generate(ctx context.Context, caps Caps, temperature floa c.compactLogLocked() // le tour est fini : coalesce ses tokens pour garder le journal petit c.mu.Unlock() c.persist() - // Notification Web Push : ce chemin (generate) ne sert QUE les tours + // File d'attente : interrompu par un stop, l'utilisateur reprend la main + // et les messages en file sont abandonnés. Sinon on démarre le prochain + // tour depuis la file — ceux qui n'ont pas pu être injectés en cours de + // route sont traités, dans l'ordre. + if ctx.Err() != nil { + c.dropQueued(epoch) + return + } + // Notification Web Push, seulement quand plus rien ne suit : un message + // en file qui démarre aussitôt n'est pas une « réponse prête », et sa + // propre fin notifiera. Ce chemin (generate) ne sert QUE les tours // utilisateur — les tâches de fond passent par RunAutonomous et ont leur // propre notification — donc pas de doublon. Détaché : l'envoi HTTP vers // le service de push ne doit pas retenir la fin du tour. Corps générique - // (pas d'extrait de réponse) : la notif transite par Apple/Google. - // - // ctx.Err() != nil = tour interrompu par un « stop » : pas de notification, - // l'utilisateur est là et a coupé volontairement. - if hasPushSubs() && ctx.Err() == nil { + // (pas d'extrait de réponse) : la notif transite par Apple/Google. Pas + // de notification après un « stop » (retour ci-dessus) : l'utilisateur + // est là et a coupé volontairement. + if !c.startQueuedIfAny() && hasPushSubs() { go sendPushToAll("Loki", "Réponse prête · "+fmtDurFR(time.Since(turnStart))) } }() @@ -680,7 +836,7 @@ func (c *Conversation) generate(ctx context.Context, caps Caps, temperature floa c.appendDelta(epoch, map[string]any{"content": ev.Content}) } return true // génération détachée : on ne s'interrompt jamais sur un abonné - }) + }, func() []Message { return c.drainQueued(epoch) }) // 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 @@ -821,6 +977,7 @@ func (c *Conversation) Reset() { c.epoch++ c.Generating = false c.cancel = nil + c.queued = nil // les messages en attente visaient le fil qu'on vient de vider c.cond.Broadcast() c.mu.Unlock() c.persist() @@ -1028,14 +1185,14 @@ func historyHead(log []LogEvent, cut, hidden int) []map[string]any { // Subscribe : abonnement sans pagination (rejeu complet). func (c *Conversation) Subscribe(ctx context.Context, from int, emit func(map[string]any) bool) { - c.SubscribeTail(ctx, from, -1, emit) + c.SubscribeTail(ctx, from, -1, "", emit) } // SubscribeTail diffuse les événements au client via emit : d'abord un REPLAY coalescé // de Log[from:] (léger), puis un caught_up, puis le DIRECT événement par événement. // 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). -func (c *Conversation) SubscribeTail(ctx context.Context, from int, tail int, emit func(map[string]any) bool) { +func (c *Conversation) SubscribeTail(ctx context.Context, from int, tail int, convID string, emit func(map[string]any) bool) { // Réveille les attentes de cond quand la connexion se ferme. go func() { <-ctx.Done() @@ -1067,8 +1224,25 @@ func (c *Conversation) SubscribeTail(ctx context.Context, from int, tail int, em // la conversation ouverte apparaît tronquée, voire vide. En session normale, le // client ne peut jamais avoir un from SUPÉRIEUR au dernier Seq du journal ; s'il // l'est, son curseur vient d'une autre session → on repart du début. + // + // ⚠️ Même famille, cas plus sournois (AJEAN 0.14.0) : un AUTRE appareil a + // changé de discussion pendant que ce client était déconnecté, et son `from` + // est PLUS BAS que le dernier Seq de la nouvelle. Le rejeu greffait alors la + // fin de la nouvelle discussion sur le début de l'ancienne restée à l'écran — + // deux fils fusionnés jusqu'au prochain rafraîchissement. Le client renvoie + // l'id de la discussion qu'il affiche (convID) : s'il ne correspond plus, on + // lui ordonne d'abord de vider l'écran (reset), puis on rejoue tout. + curID := getStr(bkChat, ckActive) + staleConv := from > 0 && convID != "" && convID != curID if n := len(snapshot); n > 0 && from > snapshot[n-1].Seq { from = 0 + staleConv = staleConv || convID != "" + } + if staleConv { + from = 0 + if !emit(map[string]any{"reset": true, "id": curID}) { + return + } } // Chargement initial d'un client paginé : on ne rejoue que la fin du fil, // précédée de {history_more: N}. Sur un long fil d'agent, tout rejouer @@ -1096,7 +1270,7 @@ func (c *Conversation) SubscribeTail(ctx context.Context, from int, tail int, em last = s } } - if !emit(map[string]any{"caught_up": true}) { + if !emit(map[string]any{"caught_up": true, "id": curID}) { return } // Taille du contexte connue MAINTENANT : le journal ne rejoue que les @@ -1141,7 +1315,8 @@ func (c *Conversation) SubscribeTail(ctx context.Context, from int, tail int, em } } c.mu.Unlock() - if !emit(map[string]any{"reset": true}) { + // id : la discussion désormais affichée (lue hors verrou — base). + if !emit(map[string]any{"reset": true, "id": getStr(bkChat, ckActive)}) { return } for _, ev := range head { diff --git a/internal/loki/chat_history_tail_test.go b/internal/loki/chat_history_tail_test.go index bf246ab..f592492 100644 --- a/internal/loki/chat_history_tail_test.go +++ b/internal/loki/chat_history_tail_test.go @@ -56,7 +56,7 @@ func TestSubscribeTailPaginates(t *testing.T) { collect := func(tail int) (more any, users int) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() - c.SubscribeTail(ctx, 0, tail, func(m map[string]any) bool { + c.SubscribeTail(ctx, 0, tail, "", func(m map[string]any) bool { if v, ok := m["history_more"]; ok { more = v } diff --git a/internal/loki/chat_queue_test.go b/internal/loki/chat_queue_test.go new file mode 100644 index 0000000..4d058c9 --- /dev/null +++ b/internal/loki/chat_queue_test.go @@ -0,0 +1,82 @@ +package loki + +import ( + "context" + "testing" +) + +// Pendant une génération de chat, un envoi est mis en file ; il ressort à la +// frontière d'étape suivante, journalisé comme une bulle utilisateur. +func TestQueueDuringGeneration(t *testing.T) { + c := newTestConv() + c.Generating = true + queued, err := c.EnqueueOrStart("cid-1", "précision", nil, Caps{}, 0.7) + if err != nil || !queued { + t.Fatalf("mise en file attendue, got queued=%v err=%v", queued, err) + } + // Réessai du même envoi (réponse perdue) : pas de doublon. + if _, err := c.EnqueueOrStart("cid-1", "précision", nil, Caps{}, 0.7); err != ErrDupSend { + t.Fatalf("doublon attendu, got %v", err) + } + msgs := c.drainQueued(c.epoch) + if len(msgs) != 1 || msgs[0].Role != "user" { + t.Fatalf("un message user attendu, got %#v", msgs) + } + if len(c.Log) != 1 || c.Log[0].Delta["user"] != "précision" { + t.Fatalf("le message injecté doit apparaître dans le fil, journal : %#v", c.Log) + } + if again := c.drainQueued(c.epoch); len(again) != 0 { + t.Fatal("la file doit être vidée par drainQueued") + } +} + +// Une tâche planifiée qui occupe le modèle garde le refus : sa fin ne dépile rien. +func TestQueueRefusedDuringTask(t *testing.T) { + c := newTestConv() + c.Generating = true + c.runningTaskName = "veille" + if queued, _ := c.EnqueueOrStart("", "x", nil, Caps{}, 0.7); queued { + t.Fatal("pas de file pendant une tâche planifiée") + } +} + +// Stop : la file est abandonnée et le client en est prévenu. +func TestDropQueuedAnnounces(t *testing.T) { + c := newTestConv() + c.Generating = true + _, _ = c.EnqueueOrStart("", "a", nil, Caps{}, 0.7) + c.dropQueued(c.epoch) + if len(c.queued) != 0 || len(c.Log) != 1 || c.Log[0].Delta["queue_dropped"] != 1 { + t.Fatalf("abandon non signalé : file=%d journal=%#v", len(c.queued), c.Log) + } +} + +// Un client dont la discussion affichée n'est plus l'active reçoit d'abord un +// reset, pour ne pas greffer le nouveau fil sur l'ancien. +func TestSubscribeStaleConvResets(t *testing.T) { + testHome(t) + _ = putStr(bkChat, ckActive, "c-nouvelle") + c := newTestConv() + for i := 0; i < 3; i++ { + c.appendDelta(c.epoch, map[string]any{"user": "q"}) + } + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + var first map[string]any + c.SubscribeTail(ctx, 2, -1, "c-ancienne", func(m map[string]any) bool { + if _, pad := m["pad"]; pad { + return true + } + if first == nil { + first = m + } + if _, ok := m["caught_up"]; ok { + cancel() + return false + } + return true + }) + if first["reset"] != true || first["id"] != "c-nouvelle" { + t.Fatalf("reset attendu en premier, got %v", first) + } +} diff --git a/internal/loki/llm_client.go b/internal/loki/llm_client.go index e130291..9bc789d 100644 --- a/internal/loki/llm_client.go +++ b/internal/loki/llm_client.go @@ -872,7 +872,12 @@ func toolCallLabel(name string, args map[string]any) string { return label } -func runChat(ctx context.Context, messages []Message, temperature float64, caps Caps, cb ChatCallback) ([]Message, error) { +// injectQueued (optionnel, variadique pour ne pas toucher aux appels hors chat) +// est consulté à CHAQUE frontière d'étape de la boucle d'outils : il renvoie les +// messages utilisateur mis en file PENDANT la génération, pour que le modèle les +// prenne en compte dans la SUITE de sa réponse (AJEAN 0.14.0) au lieu +// d'obliger à arrêter puis relancer. Ajoutés à `messages` ET à `extra`. +func runChat(ctx context.Context, messages []Message, temperature float64, caps Caps, cb ChatCallback, injectQueued ...func() []Message) ([]Message, error) { var extra []Message tools := EnabledTools(caps) // Some backends (vanilla llama.cpp builds) don't populate `reasoning_content` @@ -945,6 +950,20 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps // qui tourne en rond n'avait rien en face de lui sauf le bouton stop. toolRuns, budgetNudges, budget := 0, 0, agentBudget() for iter := 0; ; iter++ { + // Ajout en cours de réponse : entre deux étapes (après un appel d'outil ou + // une relance), on injecte les messages mis en file par l'utilisateur. Pas + // à iter==0 : le message initial du tour est déjà dans `messages`. Ni après + // un stop : la boucle repasse ici une fois l'outil interrompu, et le + // message en file serait journalisé alors que l'utilisateur a repris la + // main (dropQueued l'abandonne à la fin du tour). + if iter > 0 && ctx.Err() == nil { + for _, inject := range injectQueued { + if q := inject(); len(q) > 0 { + messages = append(messages, q...) + extra = append(extra, q...) + } + } + } // Rappel injecté EN FIN d'historique : le préfixe déjà en cache côté // llama-server reste valide, seul le nouveau message est à traiter. if msg := budgetNudge(toolRuns, budget, budgetNudges); msg != "" { diff --git a/internal/loki/ui/index.html b/internal/loki/ui/index.html index 70eb78b..a473aad 100644 --- a/internal/loki/ui/index.html +++ b/internal/loki/ui/index.html @@ -1245,6 +1245,19 @@ html[data-hide-stats="1"] #chat .msg .statline:not(.has-work){display:none} #attach-list::-webkit-scrollbar{display:none} #attach-list.show{display:flex} #composer:has(#attach-list.show){--fade:6px} +/* File d'attente : messages envoyés pendant une réponse, pas encore pris en + compte. Même bande que les pièces jointes, en gris. */ +#queue-list{order:0;display:none;flex-direction:column;align-items:flex-end;gap:4px;margin:0 var(--col);padding:0 2px 8px} +#queue-list.show{display:flex} +#composer:has(#queue-list.show){--fade:6px} +#queue-list .queued-msg{max-width:80%;padding:5px 12px;border-radius:12px;background:var(--panel); + border:1px dashed var(--border);color:var(--dim);font-size:12px;line-height:1.35; + white-space:pre-wrap;overflow:hidden;text-overflow:ellipsis;display:-webkit-box;-webkit-line-clamp:3;-webkit-box-orient:vertical; + box-shadow:0 2px 10px rgba(0,0,0,.10)} +#queue-list .queued-msg::before{content:"en attente · ";font-style:italic;opacity:.8} +/* Pendant une réponse, envoyer reste possible dès qu'il y a du texte : le + message part en file (voir send()). Stop reste à côté. */ +html[data-busy="1"][data-hastext="1"] #send{display:inline-flex} /* Posées sur le fil (et non sur la carte) : il leur faut un fond opaque et un cadre, sinon le texte de la conversation transparaît derrière. */ #attach-list .chip-file{flex:0 0 auto;background:var(--panel);border-color:var(--border); @@ -2372,6 +2385,10 @@ html[data-files="1"] #files-btn{color:var(--accent)} + +
@@ -6944,6 +6961,9 @@ function onKey(e){ if(e.key==='Enter' && !e.shiftKey){ e.preventDefault(); send( function autoGrow(ta){ ta = ta || document.getElementById('input'); if(!ta) return; + // Du texte à envoyer ? Pendant une réponse, c'est ce qui fait réapparaître + // « envoyer » à côté de stop (message mis en file, voir send()). + if(ta.id==='input') document.documentElement.setAttribute('data-hastext', ta.value.trim() ? '1' : '0'); ta.style.height='auto'; const max=parseInt(getComputedStyle(ta).maxHeight,10)||200; ta.style.height=Math.min(ta.scrollHeight, max)+'px'; @@ -6962,6 +6982,32 @@ let lastSeq=0, streamAbort=null; // caught_up). const HIST_TAIL=20; let HIST_FULL=false, HIST_RESTORE=null; +// Discussion AFFICHÉE (id reçu au caught_up/reset). Renvoyée à chaque +// (ré)abonnement : si un autre appareil a changé de discussion pendant une +// coupure, le serveur ordonne un reset au lieu de greffer le nouveau fil sur +// l'ancien (AJEAN 0.14.0). +let CONV_ID=''; +// File d'attente (AJEAN 0.14.0) : messages envoyés PENDANT une réponse, montrés +// en gris au-dessus de la carte jusqu'à ce que le flux les confirme (delta +// `user` du même texte) — injectés en cours de réponse ou au tour suivant. +function queueAdd(text){ + const box=document.getElementById('queue-list'); if(!box) return null; + const el=document.createElement('div'); el.className='queued-msg'; el.textContent=text||'(pièce jointe)'; + el.dataset.text=text; + box.appendChild(el); box.classList.add('show'); + return el; +} +function queueSync(){ const box=document.getElementById('queue-list'); if(box) box.classList.toggle('show', !!box.children.length); } +function queueTake(text){ + const box=document.getElementById('queue-list'); if(!box) return; + const el=[...box.children].find(e=>e.dataset.text===text); + if(el) el.remove(); + queueSync(); +} +function queueClear(){ const box=document.getElementById('queue-list'); if(box) box.textContent=''; queueSync(); } +// Identifiant d'envoi, stable d'un réessai à l'autre : le serveur ne met pas +// deux fois le même message en file si la réponse s'est perdue en route. +function newCID(){ try{ return crypto.randomUUID(); }catch(_){ return Date.now().toString(36)+Math.random().toString(36).slice(2); } } // Bulle « en attente » : affichée EN GRIS dès l'appui sur envoyer, avant tout // aller-retour réseau. Le message ne disparaît donc plus de l'écran entre la // frappe et la réponse du serveur. Elle s'éclaircit (classe retirée) quand @@ -7370,6 +7416,7 @@ function handleDelta(d){ // pendant le rejeu, où les `ts` sont vieux de plusieurs heures. if(!REPLAYING && typeof d.ts==='number' && d.ts>0) TS_SKEW = Date.now() - d.ts; if(d.caught_up){ + if(typeof d.id==='string') CONV_ID=d.id; settleBlocks(); // rendre le dernier bloc rejoué à sa valeur exacte // Fin du replay initial : on saute en bas puis on révèle (une seule fois — pas // sur les reconnexions, pour ne pas te ramener en bas si tu lisais plus haut). @@ -7396,7 +7443,8 @@ function handleDelta(d){ // a pu être déclenchée depuis un autre appareil. Les fichiers suivent : ils // appartiennent à la discussion, le panneau doit changer avec elle. if(d.history_more!==undefined){ showHistoryMore(d.history_more); return; } - if(d.reset!==undefined){ HIST_FULL=false; smoothReset(); cancelRender(); stopWorkTimer(); resetConvSpeed(); const tc=document.getElementById('turn-clock'); if(tc) tc.remove(); PENDING=null; document.getElementById('chat').innerHTML=''; newTurn(); setCtxUsed(0); lastSeq=0; setBusy(false); if(typeof loadConversations==='function') loadConversations(); if(typeof filesOnConvChange==='function') filesOnConvChange(); if(typeof modeOnConvChange==='function') modeOnConvChange(); return; } + if(d.queue_dropped!==undefined){ if(!REPLAYING){ queueClear(); toast('messages en attente abandonnés'); } return; } + if(d.reset!==undefined){ if(typeof d.id==='string') CONV_ID=d.id; queueClear(); HIST_FULL=false; smoothReset(); cancelRender(); stopWorkTimer(); resetConvSpeed(); const tc=document.getElementById('turn-clock'); if(tc) tc.remove(); PENDING=null; document.getElementById('chat').innerHTML=''; newTurn(); setCtxUsed(0); lastSeq=0; setBusy(false); if(typeof loadConversations==='function') loadConversations(); if(typeof filesOnConvChange==='function') filesOnConvChange(); if(typeof modeOnConvChange==='function') modeOnConvChange(); return; } // --- Mode code (20-mode.js) ----------------------------------------------- if(d.mode!==undefined){ if(typeof applyModeDelta==='function') applyModeDelta(d.mode); return; } if(d.code_hint){ if(typeof showCodeHint==='function') showCodeHint(); return; } @@ -7413,6 +7461,7 @@ function handleDelta(d){ // donc présent au rejeu comme en direct. Repris par tagModel sur chaque // bulle de réponse du tour. T.model = d.model || ''; + if(!REPLAYING) queueTake(d.user); // message en file désormais pris en compte let el=PENDING; if(!confirmPending(d.user)) el=addMsg('user', d.user); // Pièces jointes du tour : rendues DANS la bulle. La bulle en attente en @@ -7524,7 +7573,7 @@ async function connectStream(){ while(document.hidden){ await new Promise(res=>setTimeout(res, 500)); } streamAbort=new AbortController(); try{ - const r=await jfetch('/api/chat',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({from:lastSeq,tail:HIST_FULL?0:HIST_TAIL}),signal:streamAbort.signal}); + const r=await jfetch('/api/chat',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({from:lastSeq,tail:HIST_FULL?0:HIST_TAIL,conv_id:CONV_ID}),signal:streamAbort.signal}); if(REPLAYING) setChatLoading('chargement de la conversation…'); const reader=r.body.getReader(); const dec=new TextDecoder(); let buf=''; while(true){ @@ -7594,17 +7643,23 @@ function loadFullHistory(){ // erreur qu'après plusieurs échecs ET vérification que rien ne tourne — plus de // « network error » alarmiste alors que l'IA répond quand même. async function send(){ - if(busy) return; // Garde-fou : le bouton est déjà désactivé, mais l'Entrée passe aussi par ici. if(STATUS_SEEN && !MODEL_READY){ toast('le modèle n\'est pas encore prêt'); return; } const ta=document.getElementById('input'); const text=ta.value.trim(); // Un envoi sans texte est légitime s'il porte une pièce jointe (« tiens, regarde »). if(!text && !ATTACH.length) return; ta.value=''; autoGrow(ta); - // Le message s'affiche TOUT DE SUITE, en gris : il ne disparaît plus le temps - // de l'aller-retour. Il s'éclaircit quand le flux le confirme (confirmPending). - addPending(text); - const fail=(m)=>{ clearPending(); toast(m); ta.value=text; autoGrow(ta); }; + // Réponse en cours : le message part EN FILE (AJEAN 0.14.0) — il sera pris en + // compte à la prochaine étape de la réponse, ou au tour suivant. Il s'affiche + // en attente au-dessus de la carte, pas dans le fil qui s'écrit encore. + const queuedAtSend = busy; + let qel=null; + // Sinon le message s'affiche TOUT DE SUITE dans le fil, en gris : il ne + // disparaît plus le temps de l'aller-retour. Il s'éclaircit quand le flux le + // confirme (confirmPending). + if(queuedAtSend) qel=queueAdd(text); else addPending(text); + const fail=(m)=>{ if(qel){ qel.remove(); queueSync(); } else clearPending(); toast(m); ta.value=text; autoGrow(ta); }; + const cid=newCID(); // C'est ici que les fichiers partent vers le serveur — pas avant. Les pastilles // ne sont retirées qu'une fois le message accepté : tant qu'il n'est pas parti, // on doit pouvoir en enlever une, et un échec doit rester visible. @@ -7612,13 +7667,20 @@ async function send(){ if(!text && !files.length){ fail('aucun fichier n\'a pu être déposé'); return; } // Les pastilles passent dans la bulle en attente : le message porte ses // fichiers dès l'envoi, sans attendre l'aller-retour. - if(PENDING) addMsgFiles(PENDING, attachSent()); + if(PENDING && !queuedAtSend) addMsgFiles(PENDING, attachSent()); for(let attempt=0; attempt<3; attempt++){ try{ - const r=await jfetch('/api/chat/send',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({message:text,files:files,ctx_used:CTX_USED})}); - if(r.status===409 || r.ok) clearAttach(); - if(r.status===409) return; // déjà en cours (notre envoi a abouti) → OK - if(r.ok) return; // la bulle + les tokens arrivent par le flux + const r=await jfetch('/api/chat/send',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({message:text,files:files,ctx_used:CTX_USED,cid})}); + if(r.ok) clearAttach(); + // 409 = une tâche planifiée occupe le modèle : pas de file dans ce cas. + if(r.status===409){ let m='modèle occupé'; try{ m=(await r.json()).error||m; }catch(_){} fail(m); return; } + if(r.ok){ + // Parti tout de suite alors qu'on le croyait en file (le tour venait de + // finir) : la pastille disparaît, la bulle arrive par le flux. + let j={}; try{ j=await r.json(); }catch(_){} + if(qel && j.queued===false){ qel.remove(); queueSync(); } + return; // la bulle + les tokens arrivent par le flux + } if(r.status<500){ let m='erreur'; try{ m=(await r.json()).error||m; }catch(_){} fail(m); return; } }catch(e){ /* réseau : on retente */ } await new Promise(res=>setTimeout(res, 600)); diff --git a/internal/loki/ui/src/index.tmpl.html b/internal/loki/ui/src/index.tmpl.html index dd4a2e0..e1574b1 100644 --- a/internal/loki/ui/src/index.tmpl.html +++ b/internal/loki/ui/src/index.tmpl.html @@ -163,6 +163,10 @@ document.documentElement.setAttribute('data-side',localStorage.getItem('loki-sid + +
diff --git a/internal/loki/ui/src/js/08-chat-render.js b/internal/loki/ui/src/js/08-chat-render.js index f649aa3..2d0c908 100644 --- a/internal/loki/ui/src/js/08-chat-render.js +++ b/internal/loki/ui/src/js/08-chat-render.js @@ -572,6 +572,9 @@ function onKey(e){ if(e.key==='Enter' && !e.shiftKey){ e.preventDefault(); send( function autoGrow(ta){ ta = ta || document.getElementById('input'); if(!ta) return; + // Du texte à envoyer ? Pendant une réponse, c'est ce qui fait réapparaître + // « envoyer » à côté de stop (message mis en file, voir send()). + if(ta.id==='input') document.documentElement.setAttribute('data-hastext', ta.value.trim() ? '1' : '0'); ta.style.height='auto'; const max=parseInt(getComputedStyle(ta).maxHeight,10)||200; ta.style.height=Math.min(ta.scrollHeight, max)+'px'; diff --git a/internal/loki/ui/src/js/09-stream.js b/internal/loki/ui/src/js/09-stream.js index 1deb730..695e109 100644 --- a/internal/loki/ui/src/js/09-stream.js +++ b/internal/loki/ui/src/js/09-stream.js @@ -11,6 +11,32 @@ let lastSeq=0, streamAbort=null; // caught_up). const HIST_TAIL=20; let HIST_FULL=false, HIST_RESTORE=null; +// Discussion AFFICHÉE (id reçu au caught_up/reset). Renvoyée à chaque +// (ré)abonnement : si un autre appareil a changé de discussion pendant une +// coupure, le serveur ordonne un reset au lieu de greffer le nouveau fil sur +// l'ancien (AJEAN 0.14.0). +let CONV_ID=''; +// File d'attente (AJEAN 0.14.0) : messages envoyés PENDANT une réponse, montrés +// en gris au-dessus de la carte jusqu'à ce que le flux les confirme (delta +// `user` du même texte) — injectés en cours de réponse ou au tour suivant. +function queueAdd(text){ + const box=document.getElementById('queue-list'); if(!box) return null; + const el=document.createElement('div'); el.className='queued-msg'; el.textContent=text||'(pièce jointe)'; + el.dataset.text=text; + box.appendChild(el); box.classList.add('show'); + return el; +} +function queueSync(){ const box=document.getElementById('queue-list'); if(box) box.classList.toggle('show', !!box.children.length); } +function queueTake(text){ + const box=document.getElementById('queue-list'); if(!box) return; + const el=[...box.children].find(e=>e.dataset.text===text); + if(el) el.remove(); + queueSync(); +} +function queueClear(){ const box=document.getElementById('queue-list'); if(box) box.textContent=''; queueSync(); } +// Identifiant d'envoi, stable d'un réessai à l'autre : le serveur ne met pas +// deux fois le même message en file si la réponse s'est perdue en route. +function newCID(){ try{ return crypto.randomUUID(); }catch(_){ return Date.now().toString(36)+Math.random().toString(36).slice(2); } } // Bulle « en attente » : affichée EN GRIS dès l'appui sur envoyer, avant tout // aller-retour réseau. Le message ne disparaît donc plus de l'écran entre la // frappe et la réponse du serveur. Elle s'éclaircit (classe retirée) quand @@ -419,6 +445,7 @@ function handleDelta(d){ // pendant le rejeu, où les `ts` sont vieux de plusieurs heures. if(!REPLAYING && typeof d.ts==='number' && d.ts>0) TS_SKEW = Date.now() - d.ts; if(d.caught_up){ + if(typeof d.id==='string') CONV_ID=d.id; settleBlocks(); // rendre le dernier bloc rejoué à sa valeur exacte // Fin du replay initial : on saute en bas puis on révèle (une seule fois — pas // sur les reconnexions, pour ne pas te ramener en bas si tu lisais plus haut). @@ -445,7 +472,8 @@ function handleDelta(d){ // a pu être déclenchée depuis un autre appareil. Les fichiers suivent : ils // appartiennent à la discussion, le panneau doit changer avec elle. if(d.history_more!==undefined){ showHistoryMore(d.history_more); return; } - if(d.reset!==undefined){ HIST_FULL=false; smoothReset(); cancelRender(); stopWorkTimer(); resetConvSpeed(); const tc=document.getElementById('turn-clock'); if(tc) tc.remove(); PENDING=null; document.getElementById('chat').innerHTML=''; newTurn(); setCtxUsed(0); lastSeq=0; setBusy(false); if(typeof loadConversations==='function') loadConversations(); if(typeof filesOnConvChange==='function') filesOnConvChange(); if(typeof modeOnConvChange==='function') modeOnConvChange(); return; } + if(d.queue_dropped!==undefined){ if(!REPLAYING){ queueClear(); toast('messages en attente abandonnés'); } return; } + if(d.reset!==undefined){ if(typeof d.id==='string') CONV_ID=d.id; queueClear(); HIST_FULL=false; smoothReset(); cancelRender(); stopWorkTimer(); resetConvSpeed(); const tc=document.getElementById('turn-clock'); if(tc) tc.remove(); PENDING=null; document.getElementById('chat').innerHTML=''; newTurn(); setCtxUsed(0); lastSeq=0; setBusy(false); if(typeof loadConversations==='function') loadConversations(); if(typeof filesOnConvChange==='function') filesOnConvChange(); if(typeof modeOnConvChange==='function') modeOnConvChange(); return; } // --- Mode code (20-mode.js) ----------------------------------------------- if(d.mode!==undefined){ if(typeof applyModeDelta==='function') applyModeDelta(d.mode); return; } if(d.code_hint){ if(typeof showCodeHint==='function') showCodeHint(); return; } @@ -462,6 +490,7 @@ function handleDelta(d){ // donc présent au rejeu comme en direct. Repris par tagModel sur chaque // bulle de réponse du tour. T.model = d.model || ''; + if(!REPLAYING) queueTake(d.user); // message en file désormais pris en compte let el=PENDING; if(!confirmPending(d.user)) el=addMsg('user', d.user); // Pièces jointes du tour : rendues DANS la bulle. La bulle en attente en @@ -573,7 +602,7 @@ async function connectStream(){ while(document.hidden){ await new Promise(res=>setTimeout(res, 500)); } streamAbort=new AbortController(); try{ - const r=await jfetch('/api/chat',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({from:lastSeq,tail:HIST_FULL?0:HIST_TAIL}),signal:streamAbort.signal}); + const r=await jfetch('/api/chat',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({from:lastSeq,tail:HIST_FULL?0:HIST_TAIL,conv_id:CONV_ID}),signal:streamAbort.signal}); if(REPLAYING) setChatLoading('chargement de la conversation…'); const reader=r.body.getReader(); const dec=new TextDecoder(); let buf=''; while(true){ @@ -643,17 +672,23 @@ function loadFullHistory(){ // erreur qu'après plusieurs échecs ET vérification que rien ne tourne — plus de // « network error » alarmiste alors que l'IA répond quand même. async function send(){ - if(busy) return; // Garde-fou : le bouton est déjà désactivé, mais l'Entrée passe aussi par ici. if(STATUS_SEEN && !MODEL_READY){ toast('le modèle n\'est pas encore prêt'); return; } const ta=document.getElementById('input'); const text=ta.value.trim(); // Un envoi sans texte est légitime s'il porte une pièce jointe (« tiens, regarde »). if(!text && !ATTACH.length) return; ta.value=''; autoGrow(ta); - // Le message s'affiche TOUT DE SUITE, en gris : il ne disparaît plus le temps - // de l'aller-retour. Il s'éclaircit quand le flux le confirme (confirmPending). - addPending(text); - const fail=(m)=>{ clearPending(); toast(m); ta.value=text; autoGrow(ta); }; + // Réponse en cours : le message part EN FILE (AJEAN 0.14.0) — il sera pris en + // compte à la prochaine étape de la réponse, ou au tour suivant. Il s'affiche + // en attente au-dessus de la carte, pas dans le fil qui s'écrit encore. + const queuedAtSend = busy; + let qel=null; + // Sinon le message s'affiche TOUT DE SUITE dans le fil, en gris : il ne + // disparaît plus le temps de l'aller-retour. Il s'éclaircit quand le flux le + // confirme (confirmPending). + if(queuedAtSend) qel=queueAdd(text); else addPending(text); + const fail=(m)=>{ if(qel){ qel.remove(); queueSync(); } else clearPending(); toast(m); ta.value=text; autoGrow(ta); }; + const cid=newCID(); // C'est ici que les fichiers partent vers le serveur — pas avant. Les pastilles // ne sont retirées qu'une fois le message accepté : tant qu'il n'est pas parti, // on doit pouvoir en enlever une, et un échec doit rester visible. @@ -661,13 +696,20 @@ async function send(){ if(!text && !files.length){ fail('aucun fichier n\'a pu être déposé'); return; } // Les pastilles passent dans la bulle en attente : le message porte ses // fichiers dès l'envoi, sans attendre l'aller-retour. - if(PENDING) addMsgFiles(PENDING, attachSent()); + if(PENDING && !queuedAtSend) addMsgFiles(PENDING, attachSent()); for(let attempt=0; attempt<3; attempt++){ try{ - const r=await jfetch('/api/chat/send',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({message:text,files:files,ctx_used:CTX_USED})}); - if(r.status===409 || r.ok) clearAttach(); - if(r.status===409) return; // déjà en cours (notre envoi a abouti) → OK - if(r.ok) return; // la bulle + les tokens arrivent par le flux + const r=await jfetch('/api/chat/send',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({message:text,files:files,ctx_used:CTX_USED,cid})}); + if(r.ok) clearAttach(); + // 409 = une tâche planifiée occupe le modèle : pas de file dans ce cas. + if(r.status===409){ let m='modèle occupé'; try{ m=(await r.json()).error||m; }catch(_){} fail(m); return; } + if(r.ok){ + // Parti tout de suite alors qu'on le croyait en file (le tour venait de + // finir) : la pastille disparaît, la bulle arrive par le flux. + let j={}; try{ j=await r.json(); }catch(_){} + if(qel && j.queued===false){ qel.remove(); queueSync(); } + return; // la bulle + les tokens arrivent par le flux + } if(r.status<500){ let m='erreur'; try{ m=(await r.json()).error||m; }catch(_){} fail(m); return; } }catch(e){ /* réseau : on retente */ } await new Promise(res=>setTimeout(res, 600)); diff --git a/internal/loki/ui/src/styles.css b/internal/loki/ui/src/styles.css index 540bd82..e4536fc 100644 --- a/internal/loki/ui/src/styles.css +++ b/internal/loki/ui/src/styles.css @@ -1210,6 +1210,19 @@ html[data-hide-stats="1"] #chat .msg .statline:not(.has-work){display:none} #attach-list::-webkit-scrollbar{display:none} #attach-list.show{display:flex} #composer:has(#attach-list.show){--fade:6px} +/* File d'attente : messages envoyés pendant une réponse, pas encore pris en + compte. Même bande que les pièces jointes, en gris. */ +#queue-list{order:0;display:none;flex-direction:column;align-items:flex-end;gap:4px;margin:0 var(--col);padding:0 2px 8px} +#queue-list.show{display:flex} +#composer:has(#queue-list.show){--fade:6px} +#queue-list .queued-msg{max-width:80%;padding:5px 12px;border-radius:12px;background:var(--panel); + border:1px dashed var(--border);color:var(--dim);font-size:12px;line-height:1.35; + white-space:pre-wrap;overflow:hidden;text-overflow:ellipsis;display:-webkit-box;-webkit-line-clamp:3;-webkit-box-orient:vertical; + box-shadow:0 2px 10px rgba(0,0,0,.10)} +#queue-list .queued-msg::before{content:"en attente · ";font-style:italic;opacity:.8} +/* Pendant une réponse, envoyer reste possible dès qu'il y a du texte : le + message part en file (voir send()). Stop reste à côté. */ +html[data-busy="1"][data-hastext="1"] #send{display:inline-flex} /* Posées sur le fil (et non sur la carte) : il leur faut un fond opaque et un cadre, sinon le texte de la conversation transparaît derrière. */ #attach-list .chip-file{flex:0 0 auto;background:var(--panel);border-color:var(--border); diff --git a/internal/loki/web_api.go b/internal/loki/web_api.go index dff67cb..4c78204 100644 --- a/internal/loki/web_api.go +++ b/internal/loki/web_api.go @@ -1349,6 +1349,15 @@ type chatReq struct { // chargement ; 0 = tout. nil = ancien client (app relais…), qui ne sait pas // afficher le bouton « échanges précédents » : il reçoit toujours tout. Tail *int `json:"tail"` + // ClientID = identifiant unique d'un ENVOI (/api/chat/send), stable d'un + // réessai à l'autre : le serveur ne met pas deux fois le même message en + // file (voir Conversation.recentCIDs). + ClientID string `json:"cid"` + // ConvID = id de la discussion AFFICHÉE par le client (reçu au dernier + // caught_up/reset du flux). Renvoyé à chaque (re)abonnement pour que le + // serveur détecte qu'un AUTRE appareil a changé de discussion entre-temps + // et ordonne un reset avant de rejouer — sinon deux fils fusionnaient. + ConvID string `json:"conv_id"` // Files = chemins relatifs des fichiers déposés juste avant par // /api/chat/upload ("uploads/rapport.pdf"). Ils sont annoncés au modèle en // tête du message (voir attachNote) ; le contenu, lui, reste sur le disque et diff --git a/internal/loki/web_chat.go b/internal/loki/web_chat.go index 5286ecf..3df0100 100644 --- a/internal/loki/web_chat.go +++ b/internal/loki/web_chat.go @@ -88,7 +88,7 @@ func runChatStream(ctx context.Context, body chatReq, emit func(map[string]any) if body.Tail != nil { tail = *body.Tail } - conv.SubscribeTail(ctx, body.From, tail, emit) + 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 @@ -106,8 +106,17 @@ func handleChatSend(w http.ResponseWriter, r *http.Request) { 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. + // 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 @@ -116,7 +125,7 @@ func handleChatSend(w http.ResponseWriter, r *http.Request) { sendJSON(w, code, map[string]any{"ok": false, "error": err.Error()}) return } - sendJSON(w, 200, map[string]any{"ok": true}) + sendJSON(w, 200, map[string]any{"ok": true, "queued": queued}) } // handleToolResult renvoie le résultat COMPLET d'un outil dont le flux n'a