Files
Loki/internal/loki/chat_conversation.go
T
MichaelandClaude Opus 5.5 58c2da1302 Chiffrement : les résultats d'outils et les images du fil suivent la mémoire
Les deux reprises d'AJEAN arrivées en parallèle du chiffrement écrivaient
en clair à côté de discussions chiffrées.

- Résultats complets des outils (« voir plus ») : bucket toolres ajouté à
  encryptedBuckets, écrits et relus par putStoreBytes / getStoreBytesErr.
  Mémoire verrouillée : rien n'est écrit, le flux porte le résultat entier.
  Le nettoyage des orphelins ne tourne plus quand la mémoire est
  verrouillée : l'index des discussions revenait vide et TOUT passait pour
  orphelin.
- Images du fil (chatimg/) : même enveloppe que les pages mémoire ;
  verrouillée, l'image reste en base64 dans le message. Elles n'étaient
  jamais effacées : comme un fil rechargé perd ses images (stripImageParts),
  celles de plus de 24 h partent au démarrage et à chaque bascule de
  discussion.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-02 23:10:56 +02:00

1187 lines
46 KiB
Go

package loki
import (
"context"
"encoding/json"
"fmt"
"os"
"sort"
"strings"
"sync"
"sync/atomic"
"time"
)
// É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). Depuis la
// coalescence en fin de tour (compactLogLocked), un tour terminé ne pèse plus
// qu'une poignée d'événements au lieu d'un par token → 20000 couvre des
// centaines de tours. La marge sert surtout à absorber UN tour en cours (streamé
// token par token) avant sa coalescence : un très gros tour (long raisonnement +
// réponse) ne doit pas se faire tronquer le début avant d'être compacté. 20000
// n'y suffisait pas — un seul tour à très long raisonnement tronquait déjà le
// journal de rejeu, et l'utilisateur perdait le DÉBUT de sa conversation à
// l'écran. Le journal ne pèse que de la mémoire d'affichage, jamais du contexte
// modèle : la marge est bon marché.
const maxLogEvents = 200000
// LogEvent = un événement d'affichage rejouable (un delta SSE + son numéro de
// séquence monotone + un horodatage serveur en ms). Le TS permet au client de
// calculer la vitesse (tok/s) à partir du temps RÉEL de génération — correct
// aussi bien en direct qu'au replay (où tout arrive d'un bloc côté client).
type LogEvent struct {
Seq int `json:"seq"`
TS int64 `json:"ts"`
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
// Tâche planifiée en cours (tasks_run.go). Le verrou de génération est le
// SEUL point partagé entre une tâche de fond et le chat : ces deux champs
// disent à l'interface POURQUOI elle est occupée, au lieu de laisser croire
// à une génération fantôme dans la discussion ouverte.
runningTaskID string
runningTaskName string
}
var conv = func() *Conversation {
c := &Conversation{}
c.cond = sync.NewCond(&c.mu)
return c
}()
// La conversation est persistée en base, sur la machine qui fait tourner le
// modèle. En clair : cette machine déchiffre déjà pour lancer le modèle, le
// relais reste aveugle — la persister ici ne change rien à la posture E2E.
// loadConvAttempts / loadConvRetryWait : au redémarrage du conteneur, l'ancien
// process peut encore tenir le verrou bbolt quelques centaines de ms pendant que
// le nouveau démarre. Quatre essais espacés de 250 ms couvrent ce chevauchement
// sans retarder perceptiblement un démarrage normal (où le premier essai passe).
const (
loadConvAttempts = 4
loadConvRetryWait = 250 * time.Millisecond
)
// loadConvRetries compte les ATTENTES de reprise du dernier LoadConversation
// (0 = le premier essai a suffi). Il n'existe que pour être observable.
//
// Le test de non-régression veut vérifier qu'une absence légitime — première
// installation, rien d'enregistré — ne déclenche AUCUNE reprise. Il le mesurait
// au chronomètre (« moins de 250 ms »), ce qui confond « aucune reprise » avec
// « une seule tentative, mais lente » : sur un runner CI chargé, créer et
// ouvrir la base bbolt a pris 435 ms et le test a échoué alors qu'aucune reprise
// n'avait eu lieu. On compte donc les reprises au lieu de les déduire d'un
// temps de mur, qui ne dépend pas que de nous.
var loadConvRetries atomic.Int32
// LoadConversation recharge l'état persisté au démarrage du process. Sans état
// enregistré (première fois) on part d'une conversation vide.
//
// Une lecture RATÉE n'est PAS une absence, et les confondre coûte cher ici :
// getStr rend "" aussi bien pour « clé absente » que pour « base inaccessible »,
// donc sur un verrou transitoire convEnsureActive croyait la discussion active
// inexistante, forgeait un NOUVEL identifiant et l'écrasait — le fil en cours
// devenait orphelin, en silence. On sonde donc la base avec son erreur AVANT
// toute écriture, on réessaie, et en cas d'échec durable on le dit et on ne
// touche à rien.
func LoadConversation() {
var b []byte
var err error
loadConvRetries.Store(0)
retry := func() {
loadConvRetries.Add(1)
time.Sleep(loadConvRetryWait)
}
for attempt := 0; attempt < loadConvAttempts; attempt++ {
// Sonde : lire la clé de la discussion active en gardant l'erreur. Tant
// qu'elle échoue, convEnsureActive ne doit surtout pas être appelée.
if _, err = getBytesErr(bkChat, ckActive); err != nil {
retry()
continue
}
// convEnsureActive reprend au passage le fil unique des versions
// précédentes (clé « conversation ») comme première discussion.
// getStoreBytesErr : même lecture, en DÉCHIFFRANT au besoin (mémoire
// chiffrée). Verrouillé, elle rend (nil, nil) — pas une erreur d'accès :
// on démarre alors sur un fil vide sans rien écraser, et tout revient au
// déverrouillage.
if b, err = getStoreBytesErr(bkChat, convKey(convEnsureActive())); err == nil {
break
}
retry()
}
if err != nil {
// Toujours en échec : on le DIT au lieu de repartir à vide en silence, et
// on abandonne le chargement sans rien écrire — un serveur qui refuse de
// démarrer serait pire, et l'historique sur disque reste intact.
fmt.Fprintf(os.Stderr, "[conv] base illisible au démarrage (%v) — aucune discussion chargée ; rien n'a été écrasé, l'historique est intact\n", err)
return
}
if len(b) == 0 {
return
}
conv.mu.Lock()
defer conv.mu.Unlock()
_ = json.Unmarshal(b, conv)
// Guérit les conversations d'avant le passage des captures en éphémère : un
// base64 d'image persisté était rejoué à chaque tour et dépassait le contexte.
conv.Messages = stripImageParts(conv.Messages)
go pruneChatImages()
// Une génération n'a pas pu survivre à l'arrêt du process : on repart propre.
conv.Generating = false
conv.cancel = nil
}
// loadFrom remplace l'état en mémoire par une autre discussion (b vide = fil
// neuf) et invalide les abonnés : l'epoch incrémenté leur fait vider l'écran et
// rejouer depuis zéro, exactement comme un reset. L'appelant NE doit PAS
// détenir mu.
func (c *Conversation) loadFrom(b []byte) {
c.Stop()
c.mu.Lock()
c.Messages, c.Log, c.Seq, c.CtxUsed = nil, nil, 0, 0
if len(b) > 0 {
_ = json.Unmarshal(b, c)
c.Messages = stripImageParts(c.Messages) // même guérison qu'au chargement
}
go pruneChatImages() // plus rien ne référence les images de l'ancien fil
c.Generating = false
c.cancel = nil
c.epoch++
c.cond.Broadcast()
c.mu.Unlock()
}
// reloadEncryptedStores recharge la discussion depuis le disque après un
// déverrouillage : verrouillée, la lecture au démarrage a rendu un fil vide.
// On ne recharge QUE dans ce cas — écraser un fil déjà rempli ferait perdre le
// tour en cours, qui n'a pas pu être persisté.
func reloadEncryptedStores() {
conv.mu.Lock()
empty := len(conv.Log) == 0
conv.mu.Unlock()
if empty {
LoadConversation()
}
}
// persist enregistre l'état (appelé en fin de tour et sur reset, pas à chaque
// delta) sous la discussion active, et rafraîchit ses métadonnées (titre déduit
// du premier message, date, nombre d'échanges). L'appelant NE doit PAS détenir mu.
func (c *Conversation) persist() {
c.mu.Lock()
// Images en base64 → références (chat_images.go) avant d'écrire : une photo
// jointe pesait des Mo, réécrits en entier à CHAQUE fin de tour.
refImagesInMessages(c.Messages)
b, err := json.Marshal(c)
title := convSummary(c.Messages)
turns := 0
for _, ev := range c.Log {
if _, ok := ev.Delta["user"]; ok {
turns++
}
}
c.mu.Unlock()
if err != nil {
return
}
id := convEnsureActive()
// Chiffré si la mémoire l'est et qu'elle est déverrouillée. Verrouillée,
// l'écriture est REFUSÉE plutôt que de remplacer un blob chiffré par du
// clair — le fil de ce tour reste en RAM, rien n'est perdu sur disque.
_ = putStoreBytes(bkChat, convKey(id), b)
convTouchMeta(id, title, turns)
}
// appendDelta journalise un événement d'affichage et réveille les abonnés.
// epoch est celui capturé au début du tour : si un Reset est passé entre-temps,
// l'événement appartient à l'ancienne conversation et est jeté (sinon il
// polluerait le journal tout neuf avec des Seq repartis de zéro).
func (c *Conversation) appendDelta(epoch int, delta map[string]any) {
c.mu.Lock()
if c.epoch != epoch {
c.mu.Unlock()
return
}
c.Seq++
c.Log = append(c.Log, LogEvent{Seq: c.Seq, TS: time.Now().UnixMilli(), Delta: delta})
if len(c.Log) > maxLogEvents {
c.Log = c.Log[len(c.Log)-maxLogEvents:]
}
c.cond.Broadcast()
c.mu.Unlock()
}
// evToks lit le nombre de tokens porté par un événement texte : 1 pour un delta
// brut (streaming), ou la valeur `toks` accumulée pour un événement déjà coalescé.
func evToks(d map[string]any) int {
switch v := d["toks"].(type) {
case int:
return v
case float64: // relu depuis le JSON persisté
return int(v)
}
return 1
}
// evTS0 lit l'horodatage de DÉBUT d'un événement texte (présent seulement sur les
// événements coalescés) ; sinon on retombe sur `fallback` (le TS de l'événement).
func evTS0(d map[string]any, fallback int64) int64 {
switch v := d["ts0"].(type) {
case int64:
return v
case float64:
return int64(v)
}
return fallback
}
// compactLogLocked coalesce EN PLACE les suites d'événements content /
// reasoning_content du journal d'affichage en un seul événement chacun (même
// logique que coalesceReplay). Sans ça, le journal grossit token par token et
// atteint maxLogEvents en quelques réponses → on perd le DÉBUT de la conversation
// à l'affichage. On l'appelle en fin de tour (les événements sont alors figés).
// On préserve `toks` (somme) et `ts0` (premier) pour que le compteur de vitesse
// (tok/s) reste correct au replay. Verrou détenu par l'appelant.
func (c *Conversation) compactLogLocked() {
if len(c.Log) < 2 {
return
}
textKey := func(d map[string]any) string {
if _, ok := d["content"].(string); ok {
return "content"
}
if _, ok := d["reasoning_content"].(string); ok {
return "reasoning_content"
}
return ""
}
out := make([]LogEvent, 0, len(c.Log))
var buf strings.Builder
bufKey := ""
var cur LogEvent
var toks int
var ts0 int64
var seq0 int
flush := func() {
if bufKey == "" {
return
}
// seq0 = seq du PREMIER delta fusionné (cur.Seq, lui, vaut celui du dernier).
cur.Delta = map[string]any{bufKey: buf.String(), "toks": toks, "ts0": ts0, "seq0": seq0}
out = append(out, cur)
buf.Reset()
bufKey, toks, ts0, seq0 = "", 0, 0, 0
}
for _, ev := range c.Log {
// Outils : un appel s'écrit en de nombreux événements (annonce done=false,
// frappe du corps, streaming des arguments) jusqu'au done=true, qui porte
// déjà l'état FINAL (résultat + diff). Les intermédiaires ne servent qu'à
// l'affichage EN DIRECT ; les garder dans le journal le faisait enfler sur
// une conversation agentique (gros write, appels en chaîne) jusqu'à
// maxLogEvents → troncature des plus VIEUX événements, donc perte des
// premiers messages à l'affichage. On ne conserve donc que le done=true
// (même politique que coalesceReplay au replay/export).
if tu, ok := ev.Delta["tool_used"].(map[string]any); ok {
flush()
if done, _ := tu["done"].(bool); done {
out = append(out, ev)
}
continue
}
key := textKey(ev.Delta)
if key == "" {
flush()
out = append(out, ev)
continue
}
if bufKey != "" && bufKey != key {
flush()
}
if bufKey == "" {
bufKey = key
cur = LogEvent{Seq: ev.Seq, TS: ev.TS}
ts0 = evTS0(ev.Delta, ev.TS)
seq0 = evSeq0(ev.Delta, ev.Seq)
}
buf.WriteString(ev.Delta[key].(string))
cur.Seq, cur.TS = ev.Seq, ev.TS
toks += evToks(ev.Delta)
}
flush()
c.Log = out
}
// compactAndPublish exécute UNE compaction et en publie tout le cycle de vie :
// bannière de progression, journal, installation du résultat, jauge de contexte.
// Renvoie l'historique (compacté ou inchangé) et s'il a changé.
//
// Les trois compactions (début de tour, fin de tour, bouton manuel) déroulaient
// la même douzaine de lignes recopiées, y compris le calcul du surcoût fixe :
// une correction dans l'une ne suivait pas dans les deux autres.
//
// Le surcoût, justement : le contexte RÉEL mesuré moins l'estimation des
// messages donne ce que l'estimation ne voit pas (prompt système injecté,
// schémas d'outils, gabarit de chat). On le rajoute à l'estimation d'après
// compaction, sinon la jauge s'effondre puis resaute au tour suivant.
func (c *Conversation) compactAndPublish(ctx context.Context, epoch int, phase string, msgs []Message, ctxUsed int, caps Caps) ([]Message, bool) {
c.appendDelta(epoch, map[string]any{"compacting": true})
compacted, changed := compactMessages(ctx, msgs, caps)
c.appendDelta(epoch, map[string]any{"compacting": false})
logCompact(phase, ctxUsed, msgs, compacted, changed)
if !changed {
// Seuil franchi mais compaction sans effet (torse vide, ou réduction sous
// le minimum exigé) : on le DIT. Retirer la bannière sans un mot donnait,
// vu de l'UI, « le compactage automatique ne fait rien ».
c.appendDelta(epoch, map[string]any{"compact_noop": true})
return msgs, false
}
// Les pages mémoire lues ont été résumées avec le torse : un rappel liste
// leurs noms pour que le modèle les relise au besoin (AJEAN 0.14.0). Agent
// seulement — sans lui, pas d'outil mem_read pour y donner suite.
if caps.Agent {
compacted = remindReadMemPages(compacted, msgs)
}
overhead := ctxUsed - estimateTokens(msgs)
if overhead < 0 {
overhead = 0
}
est := estimateTokens(compacted) + overhead
c.mu.Lock()
if c.epoch == epoch {
c.Messages = compacted
c.CtxUsed = est // le vrai compte reviendra avec les stats du prochain tour
}
c.mu.Unlock()
c.appendDelta(epoch, map[string]any{"compacted": true})
c.appendDelta(epoch, map[string]any{"ctx_used": est}) // fait chuter la jauge tout de suite
return compacted, true
}
// convState renvoie un instantané léger (pour /api/chat/state).
func (c *Conversation) state() map[string]any {
c.mu.Lock()
defer c.mu.Unlock()
// `turns` = nombre d'échanges, borne du curseur de portée de l'export (voir
// countTurns). Compté ici plutôt que par un appel dédié : c'est un balayage
// du journal déjà en main, et l'état est de toute façon relu à l'ouverture.
turns := 0
for _, ev := range c.Log {
if _, ok := ev.Delta["user"]; ok {
turns++
}
}
st := map[string]any{"seq": c.Seq, "generating": c.Generating, "ctx_used": c.CtxUsed, "turns": turns}
// Occupé PAR UNE TÂCHE : l'interface l'annonce (« tâche en cours ») plutôt que
// d'afficher un tour qui n'apparaîtra jamais dans le fil.
if c.runningTaskID != "" {
st["task_id"] = c.runningTaskID
st["task_name"] = c.runningTaskName
}
return st
}
// isGenerating dit si un tour est en cours. Sert à la LISTE des discussions :
// elle doit pouvoir montrer laquelle travaille dès le chargement de la page,
// sans attendre que le rejeu du flux SSE ait rattrapé son retard.
func (c *Conversation) isGenerating() bool {
c.mu.Lock()
defer c.mu.Unlock()
return c.Generating
}
// ErrBusy : une génération est déjà en cours (un seul tour à la fois).
var ErrBusy = fmt.Errorf("génération en cours")
// busyReason explique le refus à l'utilisateur. Sans ça, un message envoyé
// pendant qu'une TÂCHE planifiée tourne renvoyait « génération en cours » alors
// que le fil est vide et que rien ne bouge à l'écran — incompréhensible. On
// nomme la tâche et on dit comment reprendre la main.
func (c *Conversation) busyReason() error {
c.mu.Lock()
name := c.runningTaskName
c.mu.Unlock()
if name == "" {
return ErrBusy
}
return fmt.Errorf("la tâche planifiée « %s » occupe le modèle — attends la fin ou arrête-la (bouton stop)", name)
}
// errModelLoading est renvoyée telle quelle à l'utilisateur, dans le chat : ce
// n'est pas un défaut mais une attente, et le message doit le dire.
var errModelLoading = fmt.Errorf("⏳ Le modèle est encore en train de charger — réessaie dans quelques secondes.")
// 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.
// files = pièces jointes déjà déposées (web_upload.go). Elles sont annoncées au
// MODÈLE en tête du message, mais rendues comme pastilles dans la bulle : le fil
// doit montrer ce que l'utilisateur a écrit, pas la consigne qu'on ajoute pour lui.
func (c *Conversation) StartTurn(text string, files []attachInfo, caps Caps, temperature float64) error {
if !healthCheck() {
return errModelLoading
}
c.mu.Lock()
if c.Generating {
c.mu.Unlock()
return ErrBusy
}
c.Generating = true
ctx, cancel := context.WithCancel(context.Background())
c.cancel = cancel
// Envoi sans un mot, juste un fichier : la bulle reste vide (les pastilles
// disent tout), mais le modèle a besoin d'une demande — sans elle il reçoit
// une liste de fichiers et rien à en faire.
prompt := text
if strings.TrimSpace(prompt) == "" {
prompt = "Prends-en connaissance."
}
// Mode code : fige le rôle du tour (plan demandé → planner, sinon builder).
// La détection de bascule (puce « passer en mode Code ? ») est émise après
// la bulle utilisateur, plus bas.
if caps.Code && caps.Role == "" {
caps.Role = "builder"
if wantsPlan(text) {
caps.Role = "planner"
}
}
// Content = simple texte d'ordinaire ; format multimodal (texte + images) quand
// la vision est active et qu'une pièce jointe est une image (userMessageContent).
c.Messages = append(c.Messages, Message{Role: "user", Content: userMessageContent(files, prompt)})
epoch := c.epoch
c.mu.Unlock()
// Borne de tour + bulle utilisateur (rejouables). Persistée tout de suite :
// si le process meurt en pleine génération (crash, restart après MAJ), le
// message de l'utilisateur survit au lieu de disparaître avec le tour.
delta := map[string]any{"user": text}
// Nom du modèle qui va produire la réponse : journalisé avec la borne de
// tour, donc rejoué au chargement — la carte de réponse garde son modèle
// après un rafraîchissement, et un vieux tour garde le modèle de l'époque,
// pas celui chargé aujourd'hui.
if m := chatModelName(); m != "" {
delta["model"] = m
}
if len(files) > 0 {
delta["files"] = files
}
c.appendDelta(epoch, delta)
// En mode chat, un message qui ressemble à une tâche de code fait
// apparaître la puce « passer en mode Code ? » (une fois par discussion).
if !caps.Code {
maybeCodeHint(c, epoch, text)
}
c.persist()
if temperature == 0 {
temperature = 0.7
}
go c.generate(ctx, caps, temperature, epoch)
return nil
}
// 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 {
m := strings.TrimSpace(ReadConfig()["MODEL"])
if m == "" {
return ""
}
if i := strings.LastIndexAny(m, "/\\"); i >= 0 {
m = m[i+1:]
}
if strings.HasSuffix(strings.ToLower(m), ".gguf") {
m = m[:len(m)-len(".gguf")]
}
return m
}
// 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.
// epoch est capturé au StartTurn : si un Reset survient pendant la génération,
// tout ce que ce tour produirait ensuite (deltas, messages, persistance) est
// abandonné au lieu de ressusciter des morceaux de l'ancienne conversation.
func (c *Conversation) generate(ctx context.Context, caps Caps, temperature float64, epoch int) {
turnStart := time.Now()
defer func() {
c.mu.Lock()
stale := c.epoch != epoch
// ⚠️ Un tour périmé ne touche PAS à l'état courant. Depuis que Reset
// débloque lui-même la conversation, un tour abandonné peut se terminer
// APRÈS le démarrage du suivant : remettre Generating à false ici
// déclarerait « libre » une génération toute neuve, et l'UI afficherait
// un fil qui se remplit avec un bouton « envoyer » actif.
if !stale {
c.Generating = false
c.cancel = nil
}
c.mu.Unlock()
if stale {
return // Reset pendant le tour : Reset a déjà persisté l'état vide
}
c.appendDelta(epoch, map[string]any{"turn_done": true})
c.mu.Lock()
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
// 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 {
go sendPushToAll("Loki", "Réponse prête · "+fmtDurFR(time.Since(turnStart)))
}
}()
// 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. Le résumé est un appel modèle non streamé :
// il bloque plusieurs secondes AVANT que la vraie réponse commence, d'où la
// bannière de progression émise par compactAndPublish.
if compactWouldTrigger(msgs, ctxUsed) {
if out, changed := c.compactAndPublish(ctx, epoch, "début-tour", msgs, ctxUsed, caps); changed {
msgs = out
}
}
// Prompt système personnalisé (UI → /api/sysprompt, fichier côté serveur).
// Injecté seulement dans la vue envoyée au modèle, jamais persisté dans
// c.Messages : modifiable à chaud, effet dès le tour suivant.
// Contexte du projet actif (description, index mémoire, index des trackers),
// même traitement : injecté dans la vue envoyée, jamais persisté — voir
// projectSystemMessages.
//
// Mode agent OFF = modèle BRUT (AJEAN 0.13.13) : ni l'un ni l'autre. Sans
// agent, on veut parler au modèle nu, comme à un llama-server direct. Le
// prompt du preset et l'index mémoire décrivent des outils que ce mode n'a
// pas : le modèle se mettait à écrire des <tool_call> en clair.
final := msgs
if caps.Agent {
final = append(projectSystemMessages(), msgs...)
if sp := readSysPrompt(); sp != "" {
final = append([]Message{{Role: "system", Content: sp}}, final...)
}
}
// newBase : vue modèle publiée par une compaction survenue PENDANT le tour.
// Non-nil = elle remplace l'historique (elle contient déjà le tour en cours).
var newBase []Message
var content strings.Builder
// Le comptage EXACT du contexte vient de `usage.prompt_tokens` (option
// include_usage). Tous les moteurs ne le renvoient pas — un llama-server
// récent a cessé de le faire, et la jauge est restée bloquée à zéro sur des
// discussions de plusieurs dizaines de milliers de tokens. On note donc si
// l'usage est arrivé ; sinon on retombe sur l'estimation (celle qui pilote
// déjà la compaction), approximative mais jamais absente.
sawUsage := false
extra, _ := runChat(ctx, InjectSkills(final, caps), temperature, caps, func(ev StreamEvent) bool {
switch {
case ev.Err != nil:
c.appendDelta(epoch, map[string]any{"error": ev.Err.Error()})
case ev.ToolUsed != nil:
tu := map[string]any{
"name": ev.ToolUsed.Name, "label": ev.ToolUsed.Label,
"result": ev.ToolUsed.Result, "done": ev.ToolUsed.Done, "typing": ev.ToolUsed.Typing,
}
// Corps en cours de frappe : transitoire (l'état final est le diff), on
// ne l'ajoute que quand il est là pour ne pas gonfler chaque événement.
if ev.ToolUsed.Body != "" {
tu["body"] = ev.ToolUsed.Body
}
// Lignes +/- d'une écriture (edit / mémoire) : persistées avec le reste
// pour que le diff soit encore là après un rafraîchissement.
if len(ev.ToolUsed.Diff) > 0 {
tu["diff"] = ev.ToolUsed.Diff
tu["added"] = ev.ToolUsed.Added
tu["removed"] = ev.ToolUsed.Removed
}
if ev.ToolUsed.Image != "" {
tu["image"] = ev.ToolUsed.Image
}
// Aperçu seulement dans le flux : taille réelle + id pour « voir plus ».
if ev.ToolUsed.ResultChars > 0 {
tu["result_chars"] = ev.ToolUsed.ResultChars
}
if ev.ToolUsed.ResultID != "" {
tu["result_id"] = ev.ToolUsed.ResultID
}
c.appendDelta(epoch, map[string]any{"tool_used": tu})
case ev.NewHistory != nil:
// Compaction faite en cours de tour : elle remplace la base au lieu de
// s'ajouter à l'ancienne (voir StreamEvent.NewHistory). On retire le
// préfixe système injecté à la volée (prompt perso + skills, fusionnés en
// UN message system en tête) : il n'appartient pas à l'historique persisté
// et doit rester modifiable à chaud.
base := ev.NewHistory
for len(base) > 0 && base[0].Role == "system" {
base = base[1:]
}
newBase = append([]Message(nil), base...)
case ev.Compacting != nil:
// Compaction déclenchée pendant la boucle d'outils : même bannière que la
// compaction de début de tour.
c.appendDelta(epoch, map[string]any{"compacting": *ev.Compacting})
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 {
sawUsage = true
c.mu.Lock()
if c.epoch == epoch {
c.CtxUsed = ev.Stats.PromptTokensTotal + ev.Stats.GenTokens
}
c.mu.Unlock()
}
c.appendDelta(epoch, map[string]any{"stats": ev.Stats})
case ev.Ask != nil:
// Question structurée (outil ask) : carte à boutons dans l'UI.
c.appendDelta(epoch, map[string]any{"ask": ev.Ask})
case ev.DropReasoning:
c.appendDelta(epoch, map[string]any{"drop_reasoning": true})
case ev.Reasoning != "":
c.appendDelta(epoch, map[string]any{"reasoning_content": ev.Reasoning})
case ev.Content != "":
content.WriteString(ev.Content)
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é
})
// 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. Sauf si un Reset est passé entre-temps :
// la nouvelle conversation vide ne doit pas hériter de la fin de l'ancienne.
c.mu.Lock()
if c.epoch == epoch {
if newBase != nil {
c.Messages = newBase
}
c.Messages = append(c.Messages, extra...)
if s := content.String(); strings.TrimSpace(s) != "" {
c.Messages = append(c.Messages, Message{Role: "assistant", Content: s})
}
}
msgs = append([]Message(nil), c.Messages...)
ctxUsed = c.CtxUsed
stale := c.epoch != epoch
c.mu.Unlock()
// Moteur muet sur l'usage : on publie l'estimation, sinon la jauge reste à
// zéro et l'utilisateur ne voit jamais son contexte se remplir — ni la
// compaction arriver.
if !stale && !sawUsage {
est := estimateTokens(msgs)
c.mu.Lock()
if c.epoch == epoch {
c.CtxUsed = est
}
c.mu.Unlock()
c.appendDelta(epoch, map[string]any{"ctx_used": est})
ctxUsed = est
}
// Boucle du contrat (mode code) : vérification indépendante des critères,
// corrections, re-vérification — AVANT de rendre la main. Voir
// code_verify.go. No-op hors mode code ou sans critères.
if !stale && ctx.Err() == nil {
c.codeVerifyLoop(ctx, caps, temperature, epoch)
c.mu.Lock()
msgs = append([]Message(nil), c.Messages...)
ctxUsed = c.CtxUsed
stale = c.epoch != epoch
c.mu.Unlock()
}
// Compaction de FIN DE TOUR. C'était LE trou : le seuil n'était testé qu'au
// DÉBUT d'un tour, avec le contexte du tour précédent. Le tour qui fait
// franchir le seuil se termine donc à 81% et… rien. La jauge reste haute, seul
// le bouton manuel s'affiche, et l'utilisateur voit un compactage automatique
// qui « ne marche pas » — alors qu'il attendait simplement le message suivant.
// On compacte donc DÈS que le tour qui a franchi le seuil est fini : la jauge
// retombe tout de suite et le tour suivant démarre avec de la marge.
// Sauf si un Reset est passé (rien à compacter) ou si le tour a été annulé
// (bouton stop) : on n'enchaîne pas plusieurs secondes de résumé sur un stop.
if stale || ctx.Err() != nil || !compactWouldTrigger(msgs, ctxUsed) {
return
}
// context.Background() et non ctx : le tour est terminé, son contexte peut
// être annulé alors que cette compaction-là doit aller au bout.
c.compactAndPublish(context.Background(), epoch, "fin-tour", msgs, ctxUsed, caps)
}
// CompactNow force une compaction du contexte MAINTENANT, sans attendre le seuil
// (bouton « compacter » de l'UI). Détaché comme la génération : émet la bannière
// de progression, résume les anciens tours, remplace le torse et persiste. Les
// événements passent par le flux d'abonnement, donc tous les appareils voient la
// progression. Renvoie ErrBusy si un tour est déjà en cours.
func (c *Conversation) CompactNow() error {
if !healthCheck() {
return errModelLoading
}
c.mu.Lock()
if c.Generating {
c.mu.Unlock()
return ErrBusy
}
c.Generating = true
ctx, cancel := context.WithCancel(context.Background())
c.cancel = cancel
msgs := append([]Message(nil), c.Messages...)
lastReal := c.CtxUsed // dernier contexte réel mesuré, pour estimer le surcoût fixe
epoch := c.epoch
c.mu.Unlock()
go func() {
defer func() {
c.mu.Lock()
// Même précaution que dans generate : une compaction abandonnée par un
// Reset ne doit pas déclarer « libre » le tour qui a démarré depuis.
if c.epoch == epoch {
c.Generating = false
c.cancel = nil
}
c.mu.Unlock()
}()
if _, changed := c.compactAndPublish(ctx, epoch, "manuel", msgs, lastReal, Caps{}); changed {
c.persist()
}
}()
return nil
}
// 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.
//
// Reset DÉBLOQUE toujours, et c'est sa deuxième raison d'être. Il se contentait
// avant de vider le fil : si un tour restait coincé (moteur redémarré sous ses
// pieds, commande shell accrochée à ses tubes), Generating restait vrai pour
// toujours et tout message suivant se voyait refusé « génération en cours ».
// Vider le fil ne changeait rien, rafraîchir non plus : il fallait redémarrer le
// service. « Nouvelle conversation » est le geste qu'on tente naturellement dans
// ce cas ; il doit donc rendre la main, quoi qu'il arrive au tour abandonné, que
// le bump d'epoch réduit de toute façon au silence.
func (c *Conversation) Reset() {
c.Stop()
// Vider la discussion efface aussi ses fichiers (dépôts, captures, ce que
// l'agent y a écrit) : les messages qui les mentionnaient disparaissent, plus
// rien ne les rattacherait à quoi que ce soit. L'UI le dit avant de demander
// confirmation.
dropConvFiles(convEnsureActive())
c.mu.Lock()
c.Messages = nil
c.Log = nil
c.Seq = 0
c.CtxUsed = 0
c.epoch++
c.Generating = false
c.cancel = nil
c.cond.Broadcast()
c.mu.Unlock()
c.persist()
}
// evSeq0 lit le seq de DÉBUT d'un bloc texte : sur un événement coalescé (dont le
// Seq vaut celui du DERNIER delta fusionné), c'est le seq du PREMIER delta du bloc ;
// sur un delta brut, c'est son propre Seq.
//
// ⚠️ Sans ça : un client qui a déjà affiché une PARTIE d'un bloc (streaming en
// cours) et qui reçoit ensuite ce bloc coalescé — parce que la compaction de fin de
// tour a remplacé la suite de deltas par un seul événement de Seq plus grand que
// son dernier seq vu — le concatène à ce qu'il affichait déjà : la réponse
// apparaissait DEUX FOIS (la 1re copie tronquée à l'endroit exact où le client en
// était). Cf. `replace` ci-dessous.
func evSeq0(d map[string]any, fallback int) int {
switch v := d["seq0"].(type) {
case int:
return v
case float64: // relu depuis le JSON persisté
return int(v)
}
return fallback
}
// decorateEvent aplatit un événement du journal pour l'émission SSE et marque
// `replace` quand le client a DÉJÀ vu le début du bloc (from tombe à l'intérieur) :
// le texte envoyé est alors le bloc ENTIER, donc le client doit remplacer sa bulle
// au lieu d'y concaténer.
func decorateEvent(ev LogEvent, from int) map[string]any {
out := map[string]any{"seq": ev.Seq, "ts": ev.TS}
for k, v := range ev.Delta {
out[k] = v
}
if isTextDelta(ev.Delta) && evSeq0(ev.Delta, ev.Seq) <= from {
out["replace"] = true
}
return out
}
// isTextDelta : événement porteur de texte (content / reasoning_content).
func isTextDelta(d map[string]any) bool {
if _, ok := d["content"].(string); ok {
return true
}
_, ok := d["reasoning_content"].(string)
return ok
}
// coalesceReplay fusionne les deltas texte consécutifs (content / reasoning_content)
// d'un même bloc en UN seul événement, pour que le replay au chargement soit léger
// (quelques événements par tour au lieu de milliers de tokens). On conserve le
// nombre de tokens fusionnés (toks) et les bornes d'horodatage (ts0→ts) pour que le
// client reconstitue le compteur ET la vitesse. Les événements non-texte (user,
// tool_used, stats, turn_done…) passent tels quels.
func coalesceReplay(events []LogEvent, from int) []map[string]any {
var out []map[string]any
var buf strings.Builder
bufKey := ""
var bufSeq, bufToks, bufSeq0 int
var bufTs0, bufTs int64
flush := func() {
if bufKey == "" {
return
}
m := map[string]any{bufKey: buf.String(), "seq": bufSeq, "ts": bufTs, "ts0": bufTs0, "toks": bufToks, "seq0": bufSeq0}
// Le client a déjà affiché le début de ce bloc (from tombe dedans, ce qui
// arrive quand la compaction de fin de tour l'a fusionné) : on renvoie le
// bloc entier et on lui demande de REMPLACER sa bulle, pas d'y ajouter.
if bufSeq0 <= from {
m["replace"] = true
}
out = append(out, m)
buf.Reset()
bufKey, bufSeq, bufToks, bufTs0, bufTs, bufSeq0 = "", 0, 0, 0, 0, 0
}
// Coalescence des outils : un appel d'outil génère plein d'événements tool_used
// intermédiaires (streaming des arguments) jusqu'à un dernier avec done=true. Au
// replay, seul l'état FINAL de chaque bulle compte (les intermédiaires ne servent
// qu'à l'affichage progressif en direct). Sans ça, un long fil rejoue des milliers
// de tool_used → 2-3 s de rendu inutile sur mobile. On ne garde donc que les
// done=true, en mémorisant le dernier non-done pour le cas d'un outil interrompu.
var pendingTool map[string]any
flushTool := func() {
if pendingTool != nil {
out = append(out, pendingTool)
pendingTool = nil
}
}
for _, ev := range events {
if ev.Seq <= from {
continue
}
_, isTool := ev.Delta["tool_used"].(map[string]any)
// ⚠️ On ne « flushe » PAS l'annonce d'outil en attente sur un événement
// non-outil. Un appel d'outil s'écrit dans le journal en DEUX temps : une
// annonce (done=false, juste « je vais lancer X ») puis le résultat
// (done=true). Si un événement non-outil (stats, un bout de raisonnement)
// se glisse ENTRE les deux, flusher ici émettait l'annonce PUIS le done →
// le même outil apparaissait DEUX fois (à l'export Markdown et au replay).
// On garde donc l'annonce en attente : le done la remplace (émis une seule
// fois), et si aucun done n'arrive jamais (outil interrompu par un stop),
// le flushTool final l'émet une seule fois, en fin de fil.
// Delta texte ? (une seule clé content ou reasoning_content, valeur string)
key := ""
if s, ok := ev.Delta["content"].(string); ok {
key, _ = "content", s
} else if s, ok := ev.Delta["reasoning_content"].(string); ok {
key, _ = "reasoning_content", s
}
if key != "" {
if bufKey != "" && bufKey != key {
flush()
}
if bufKey == "" {
bufKey, bufTs0, bufSeq0 = key, evTS0(ev.Delta, ev.TS), evSeq0(ev.Delta, ev.Seq)
}
buf.WriteString(ev.Delta[key].(string))
bufSeq, bufTs = ev.Seq, ev.TS
bufToks += evToks(ev.Delta) // 1 pour un delta brut, N pour un événement déjà coalescé
continue
}
flush()
m := map[string]any{"seq": ev.Seq, "ts": ev.TS}
for k, v := range ev.Delta {
m[k] = v
}
if isTool {
tu, _ := ev.Delta["tool_used"].(map[string]any)
done, _ := tu["done"].(bool)
if done {
pendingTool = nil // les intermédiaires de cet outil sont superflus
out = append(out, m)
} else {
pendingTool = m // on retient le dernier état non terminé, sans l'émettre
}
continue
}
out = append(out, m)
}
flush()
flushTool()
return out
}
// historyTailDefault : échanges rejoués à l'ouverture d'une session (reset) pour
// un client qui gère la pagination de l'historique.
const historyTailDefault = 20
// historyCut calcule où couper le rejeu pour ne garder que les `tail` derniers
// échanges (un échange commence à un événement `user`). Renvoie le Seq à partir
// duquel rejouer (exclu) et le nombre d'échanges masqués ; (0, 0) = tout
// rejouer. Repris d'AJEAN 0.15.7.
func historyCut(log []LogEvent, tail int) (int, int) {
if tail <= 0 {
return 0, 0
}
seen := 0
for i := len(log) - 1; i >= 0; i-- {
if _, ok := log[i].Delta["user"]; !ok {
continue
}
seen++
if seen == tail {
hidden := 0
for _, ev := range log[:i] {
if _, ok := ev.Delta["user"]; ok {
hidden++
}
}
if hidden == 0 {
return 0, 0
}
return log[i].Seq - 1, hidden
}
}
return 0, 0
}
// historyHead : événements qui annoncent un historique tronqué. Les états que
// les échanges masqués avaient posés et que le client garde (mode Chat/Code,
// critères du mode Code) sont réémis tels qu'ils étaient à la coupe — sans eux,
// une discussion en mode Code rouverte s'afficherait en mode Chat.
func historyHead(log []LogEvent, cut, hidden int) []map[string]any {
out := []map[string]any{{"history_more": hidden}}
var mode, crit map[string]any
for _, ev := range log {
if ev.Seq > cut {
break
}
if _, ok := ev.Delta["mode"]; ok {
mode = ev.Delta
}
if _, ok := ev.Delta["criteria"]; ok {
crit = ev.Delta
}
}
for _, d := range []map[string]any{mode, crit} {
if d != nil {
out = append(out, d)
}
}
return out
}
// 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)
}
// 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) {
// Réveille les attentes de cond quand la connexion se ferme.
go func() {
<-ctx.Done()
c.mu.Lock()
c.cond.Broadcast()
c.mu.Unlock()
}()
// 0. Amorçage anti-buffering. Le chat E2E d'app.ajean.link traverse Cloudflare
// (ajean.link, proxy orange) : tant qu'un proxy intermédiaire n'a pas reçu assez
// d'octets, il bufferise la réponse et ne la relaie qu'en retard (symptôme :
// derniers messages qui arrivent 20-30 s après le reste, de façon intermittente
// sur Safari). Envoyer d'emblée un gros événement de padding force le proxy à
// basculer en mode streaming tout de suite. Le client ignore la clé `pad`.
if !emit(map[string]any{"pad": strings.Repeat("·", 2048)}) {
return
}
// 1. Replay coalescé (snapshot hors verrou pour ne pas bloquer la génération).
c.mu.Lock()
snapshot := append([]LogEvent(nil), c.Log...)
epoch := c.epoch
c.mu.Unlock()
// ⚠️ Garde-fou anti-« premiers messages manquants au changement de session ».
// Le client se réabonne avec from=lastSeq (le dernier Seq qu'il a vu). Les Seq
// ne sont PAS globaux : ouvrir une session PLUS ANCIENNE charge un journal dont
// les Seq sont plus bas. Si le client garde un from élevé (hérité de la session
// qu'il quittait) et se reconnecte, coalesceReplay saute tout (Seq <= from) →
// 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.
if n := len(snapshot); n > 0 && from > snapshot[n-1].Seq {
from = 0
}
// 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
// envoyait des Mo et bâtissait des milliers de bulles avant d'afficher quoi
// que ce soit.
if from == 0 && tail > 0 {
if cut, hidden := historyCut(snapshot, tail); hidden > 0 {
from = cut
for _, ev := range historyHead(snapshot, cut, hidden) {
if !emit(ev) {
return
}
}
}
}
last := from
for _, ev := range coalesceReplay(snapshot, from) {
if ctx.Err() != nil {
return
}
if !emit(ev) {
return
}
if s, ok := ev["seq"].(int); ok {
last = s
}
}
if !emit(map[string]any{"caught_up": true}) {
return
}
// Taille du contexte connue MAINTENANT : le journal ne rejoue que les
// événements du fil, et `ctx_used` n'y figure qu'aux tours où il a changé —
// souvent hors de la fenêtre rejouée. Sans cet envoi, toute page rechargée
// affichait « 0 / N (0 %) » sur une discussion pourtant bien remplie.
c.mu.Lock()
used := c.CtxUsed
c.mu.Unlock()
if used > 0 && !emit(map[string]any{"ctx_used": used}) {
return
}
// Le padding doit venir APRÈS caught_up, pas dessus. Un proxy (Cloudflare) garde
// toujours le DERNIER bout de flux en tampon jusqu'au prochain flush (~2-3 s). En
// envoyant un gros pad juste après, ce sont ses octets — et non le dernier vrai
// message — qui deviennent la « queue » qui attend : le dernier message et
// caught_up, eux, sont poussés dehors immédiatement. Le client ignore `pad`.
// 16 Ko : largement au-dessus du tampon de coalescence d'un proxy courant.
if !emit(map[string]any{"pad": strings.Repeat("·", 16384)}) {
return
}
// 2. Direct : événements granulaires au-delà de `last`.
c.mu.Lock()
for {
if ctx.Err() != nil {
c.mu.Unlock()
return
}
if c.epoch != epoch { // reset → on ordonne au client de nettoyer et on repart
epoch = c.epoch
last = 0
// Session ouverte (fil complet qui suit) chez un client paginé : on
// saute directement aux derniers échanges.
var head []map[string]any
if tail >= 0 {
if cut, hidden := historyCut(c.Log, historyTailDefault); hidden > 0 {
last = cut
// ctx_used vit dans le journal des échanges masqués : on le
// donne tel qu'il est maintenant, sinon la jauge resterait à 0.
head = append(historyHead(c.Log, cut, hidden), map[string]any{"ctx_used": c.CtxUsed})
}
}
c.mu.Unlock()
if !emit(map[string]any{"reset": true}) {
return
}
for _, ev := range head {
if !emit(ev) {
return
}
}
c.mu.Lock()
continue
}
// Copie des événements en attente SOUS verrou, émission HORS verrou : on
// n'itère jamais sur c.Log pendant que la génération peut y écrire ou que
// la troncature (maxLogEvents) peut le déplacer.
//
// c.Log est trié par Seq croissant : on trouve le premier événement neuf
// par dichotomie plutôt qu'en relisant tout. Le balayage complet coûtait
// la longueur du journal (jusqu'à 20 000) À CHAQUE TOKEN et pour CHAQUE
// appareil connecté, verrou tenu — donc au détriment de la génération
// elle-même. C'est quadratique sur un long fil.
i := sort.Search(len(c.Log), func(i int) bool { return c.Log[i].Seq > last })
var pending []LogEvent
if i < len(c.Log) {
pending = append(pending, c.Log[i:]...)
}
if len(pending) == 0 {
c.cond.Wait()
continue
}
lastEmitted := last
last = pending[len(pending)-1].Seq
c.mu.Unlock()
for _, ev := range pending {
// `lastEmitted` (avant mise à jour) sert de repère : si la compaction de
// fin de tour vient de fusionner un bloc que le client suivait en direct,
// l'événement fusionné arrive avec un Seq supérieur au sien → `replace`.
if !emit(decorateEvent(ev, lastEmitted)) {
return
}
}
c.mu.Lock()
}
}