From c3005f5b56bc71cc79e96cacb2db6fab1e09e28d Mon Sep 17 00:00:00 2001 From: Michael SCHAL Date: Sun, 4 Oct 2026 02:19:03 +0200 Subject: [PATCH] Persistance : le fsync de la discussion ne retarde plus le premier token MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit StartTurn persistait AVANT de lancer la génération : un marshal, puis deux commits bbolt (contenu, puis index), soit deux fsync et plusieurs ouvertures de base — 26-30 ms mesurés sur un NVMe Windows, bien plus sur un /mnt/user d'Unraid ou un disque de parité, payés à chaque message avant que la requête ne parte vers llama-server. Même facture au milieu d'un tour (vérification du mode Code) et pour chaque résultat d'outil coupé (« voir plus »). - chat_persist.go : un écrivain unique et ordonné. L'instantané reste pris sous c.mu (images → références, marshal, identifiant de la discussion) ; seule l'E/S part en différé. Le dernier instantané de chaque discussion l'emporte (numéro pris sous c.mu), les lots sont fusionnés en UNE transaction : contenu + index ensemble, résultats d'outils compris. - Seuls StartTurn, la persistance en cours de tour et les résultats d'outils sont asynchrones. Fin de tour, reset, bascules, compactage manuel restent synchrones et attendent tout ce qui précède ; la suppression écarte puis attend les écritures qui la visent (plus de discussion ressuscitée). - Discussion active gardée en RAM, créée sous verrou (plus de double identifiant sur base neuve) et changée sous c.mu avec le contenu : un instantané ne peut plus écrire l'ancien fil sous le nouvel identifiant. - « Voir plus » servi depuis la mémoire tant que le résultat n'est pas écrit ; verrouillage vérifié avant de rendre un id. - En plein tour, la base reçoit le journal sous sa forme de fin de tour (compactLog, version pure de compactLogLocked) ; la mémoire n'est pas touchée. - Vidage de l'écrivain avant redémarrage, « Quitter », exec et (dé)chiffrement. Le contexte vu par le modèle ne change pas : il lit c.Messages en mémoire, les résultats complets d'outils et la compaction du journal ne servent qu'à l'UI. Co-Authored-By: Claude Opus 5.5 --- internal/loki/chat_conversation.go | 69 ++-- internal/loki/chat_persist.go | 539 +++++++++++++++++++++++++++ internal/loki/chat_persist_test.go | 425 +++++++++++++++++++++ internal/loki/chat_sessions.go | 81 ++-- internal/loki/code_verify.go | 4 +- internal/loki/mem_store.go | 4 + internal/loki/store.go | 8 +- internal/loki/store_test.go | 3 + internal/loki/sys_platform_unix.go | 3 + internal/loki/sys_restart_windows.go | 1 + internal/loki/sys_tray_darwin.go | 1 + internal/loki/sys_tray_windows.go | 1 + internal/loki/tool_results.go | 22 +- internal/loki/tool_results_test.go | 1 + 14 files changed, 1082 insertions(+), 80 deletions(-) create mode 100644 internal/loki/chat_persist.go create mode 100644 internal/loki/chat_persist_test.go diff --git a/internal/loki/chat_conversation.go b/internal/loki/chat_conversation.go index a7f64e9..76dbfbe 100644 --- a/internal/loki/chat_conversation.go +++ b/internal/loki/chat_conversation.go @@ -191,11 +191,15 @@ func LoadConversation() { // 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) { +// rejouer depuis zéro, exactement comme un reset. id = discussion chargée : le +// cache de la discussion active change sous mu, avec le contenu (voir +// convActivate). L'appelant NE doit PAS détenir mu. +func (c *Conversation) loadFrom(id string, b []byte) { c.Stop() c.mu.Lock() + if c == conv { + convActiveRemember(id) + } c.Messages, c.Log, c.Seq, c.CtxUsed = nil, nil, 0, 0 c.ctxUsedLen, c.genPeak = 0, 0 c.queued = nil // file de l'ancienne discussion : elle ne suit pas la bascule @@ -224,34 +228,6 @@ func reloadEncryptedStores() { } } -// 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 @@ -303,8 +279,16 @@ func evTS0(d map[string]any, fallback int64) int64 { // 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 + c.Log = compactLog(c.Log) +} + +// compactLog est la compaction de compactLogLocked sous forme PURE : elle rend +// un nouveau journal sans toucher à celui qu'on lui passe (les événements non +// fusionnés sont partagés, pas recopiés). snapshot s'en sert pour écrire, en +// plein tour, la forme qu'aura le journal à la fin du tour. +func compactLog(log []LogEvent) []LogEvent { + if len(log) < 2 { + return log } textKey := func(d map[string]any) string { if _, ok := d["content"].(string); ok { @@ -315,7 +299,7 @@ func (c *Conversation) compactLogLocked() { } return "" } - out := make([]LogEvent, 0, len(c.Log)) + out := make([]LogEvent, 0, len(log)) var buf strings.Builder bufKey := "" var cur LogEvent @@ -332,7 +316,7 @@ func (c *Conversation) compactLogLocked() { buf.Reset() bufKey, toks, ts0, seq0 = "", 0, 0, 0 } - for _, ev := range c.Log { + for _, ev := range 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'à @@ -368,7 +352,7 @@ func (c *Conversation) compactLogLocked() { toks += evToks(ev.Delta) } flush() - c.Log = out + return out } // compactAndPublish exécute UNE compaction et en publie tout le cycle de vie : @@ -511,9 +495,10 @@ func (c *Conversation) StartTurn(text string, files []attachInfo, caps Caps, tem 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. + // Borne de tour + bulle utilisateur (rejouables). Persistée tout de suite + // (sur disque quelques ms plus tard, voir persistAsync) : 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 @@ -531,7 +516,9 @@ func (c *Conversation) StartTurn(text string, files []attachInfo, caps Caps, tem if !caps.Code { maybeCodeHint(c, epoch, text) } - c.persist() + // Instantané pris ici, écriture laissée à l'écrivain (chat_persist.go) : le + // fsync ne retarde plus le départ de la requête vers le modèle. + c.persistAsync() if temperature == 0 { temperature = 0.7 } @@ -750,7 +737,7 @@ func (c *Conversation) generate(ctx context.Context, caps Caps, temperature floa }() // Télémétrie : tout ce que ce tour envoie au moteur (étapes, compaction, // sous-agents, vérification) est rattaché à la discussion active. - ctx = withPerf(ctx, perfMain, getStr(bkChat, ckActive)) + ctx = withPerf(ctx, perfMain, convEnsureActive()) // llama-server local seulement : le preset externe garde le seul seuil, et // ses complétions ne nourrissent pas la garde de marge (compactNeeded). diff --git a/internal/loki/chat_persist.go b/internal/loki/chat_persist.go new file mode 100644 index 0000000..9dfd388 --- /dev/null +++ b/internal/loki/chat_persist.go @@ -0,0 +1,539 @@ +package loki + +// chat_persist.go — l'écriture de la discussion SORT du chemin du premier token. +// +// persist() faisait tout en ligne : marshal, puis putStoreBytes (ouverture de la +// base, commit, fsync), puis convTouchMeta (relecture de l'index, second commit, +// second fsync). Dans StartTurn, c'était AVANT `go c.generate` : 26-30 ms +// mesurés sur un NVMe Windows, bien davantage sur le /mnt/user d'Unraid (FUSE) +// ou un disque de parité — payés à chaque message avant que la requête ne parte +// vers llama-server. Même facture au milieu d'un tour (vérification du mode +// Code) et pour chaque résultat d'outil coupé (« voir plus »). +// +// Désormais : +// - l'INSTANTANÉ reste synchrone et pris sous c.mu (images → références, +// marshal, identifiant de la discussion) : rien ne peut le muter pendant +// qu'on le prend, et il dit toujours à quelle discussion il appartient ; +// - seule l'E/S bbolt part à un écrivain UNIQUE, qui fusionne ce qui s'est +// accumulé (dernier instantané de chaque discussion + résultats d'outils) en +// UNE transaction — un fsync au lieu de deux, ou de vingt sur un tour d'agent ; +// - seuls les chemins chauds sont asynchrones (StartTurn, persistance en cours +// de tour, résultats d'outils). Fin de tour, reset, bascules et suppression +// restent synchrones : ils attendent que l'écrivain ait tout vidé, eux +// compris. Aucun n'est sur le chemin du premier token, et un redémarrage +// juste après eux ne doit rien perdre. +// +// Ce qui atteint le modèle ne change pas d'un octet : il lit c.Messages en +// mémoire, et les résultats complets des outils ne servent qu'au « voir plus » +// de l'UI. Le risque accepté : un crash dans les millisecondes qui suivent un +// envoi peut perdre le message tout juste envoyé — le même que celui d'un crash +// en pleine génération, avant ce changement, pour la réponse. + +import ( + "encoding/json" + "fmt" + "os" + "sort" + "strings" + "sync" + "sync/atomic" + "time" + + bolt "go.etcd.io/bbolt" +) + +// convSnap est un instantané complet d'une discussion, prêt à écrire. +type convSnap struct { + path string // base visée, figée à la prise (voir withDBAt) + id string // discussion à laquelle appartient le contenu + seq uint64 // ordre de prise : le plus grand l'emporte toujours + body []byte // JSON de Conversation + title string // titre déduit du premier message (convSummary) + turns int + project string // projet d'une entrée d'index à créer (voir convTouchMetaIn) +} + +// toolResJob est un résultat d'outil complet à ranger dans bkToolRes. +type toolResJob struct { + path string + key string + plain string +} + +// convPersister sérialise TOUTES les écritures de discussion. Un seul écrivain à +// la fois : l'ordre des instantanés ne dépend donc que de seq, jamais d'une +// course entre deux goroutines d'écriture. +type convPersister struct { + mu sync.Mutex + cond *sync.Cond + snaps map[string]*convSnap // dernier instantané en attente, par discussion + res []toolResJob // résultats d'outils en attente, dans l'ordre + mem map[string]string // résultats pas encore écrits : « voir plus » les sert d'ici + written map[string]uint64 // seq du dernier instantané écrit, par discussion + deleted map[string]bool // discussions supprimées : plus rien ne doit les réécrire + queued uint64 // tickets délivrés + done uint64 // tickets traités (écrits, abandonnés ou en échec) + running bool +} + +var ( + persistQ = newConvPersister() + // persistSeq numérote les instantanés. Incrémenté SOUS c.mu : l'ordre des + // numéros est donc celui des états de la discussion. + persistSeq atomic.Uint64 + // persistCommit écrit un lot ; remplaçable par les tests pour observer + // l'ordre des écritures ou les ralentir. + persistCommit = commitPersistBatch +) + +func newConvPersister() *convPersister { + p := &convPersister{ + snaps: map[string]*convSnap{}, + mem: map[string]string{}, + written: map[string]uint64{}, + deleted: map[string]bool{}, + } + p.cond = sync.NewCond(&p.mu) + return p +} + +// enqueueSnap confie un instantané à l'écrivain et renvoie son ticket. Un +// instantané en attente pour la même discussion est remplacé par le plus récent : +// écrire l'ancien ne servirait qu'à l'écraser aussitôt. +func (p *convPersister) enqueueSnap(s *convSnap) uint64 { + p.mu.Lock() + defer p.mu.Unlock() + if !p.deleted[s.id] { + if cur := p.snaps[s.id]; cur == nil || s.seq > cur.seq { + p.snaps[s.id] = s + } + } + return p.kickLocked() +} + +// enqueueToolRes confie un résultat d'outil à l'écrivain. Il est servi depuis +// la mémoire tant qu'il n'est pas sur disque. +func (p *convPersister) enqueueToolRes(j toolResJob) uint64 { + p.mu.Lock() + defer p.mu.Unlock() + p.mem[j.key] = j.plain + p.res = append(p.res, j) + return p.kickLocked() +} + +// kickLocked délivre un ticket et démarre l'écrivain s'il dort. p.mu détenu. +func (p *convPersister) kickLocked() uint64 { + p.queued++ + if !p.running { + p.running = true + go p.run() + } + return p.queued +} + +// run vide la file par lots, jusqu'à ce qu'elle soit vide. +func (p *convPersister) run() { + for { + p.mu.Lock() + if len(p.snaps) == 0 && len(p.res) == 0 { + p.running = false + p.done = p.queued + p.cond.Broadcast() + p.mu.Unlock() + return + } + upto := p.queued + snaps := make([]*convSnap, 0, len(p.snaps)) + for id, s := range p.snaps { + if !p.deleted[id] && s.seq > p.written[id] { + snaps = append(snaps, s) + } + } + res := make([]toolResJob, 0, len(p.res)) + for _, j := range p.res { + if !p.deleted[toolResConv(j.key)] { + res = append(res, j) + } + } + p.snaps, p.res = map[string]*convSnap{}, nil + p.mu.Unlock() + + sort.Slice(snaps, func(i, j int) bool { return snaps[i].seq < snaps[j].seq }) + wrote, err := persistCommitByPath(snaps, res) + + p.mu.Lock() + for _, s := range wrote { + if s.seq > p.written[s.id] { + p.written[s.id] = s.seq + } + } + if err != nil { + // Les résultats restent servis depuis la mémoire : le « voir plus » de + // cette session marche encore, seul un redémarrage les perdrait. Un + // instantané raté sera remplacé par le suivant (fin de tour au plus tard). + fmt.Fprintf(os.Stderr, "[persist] écriture différée échouée (%d instantané(s), %d résultat(s) d'outil gardé(s) en mémoire) : %v\n", len(snaps)-len(wrote), len(res), err) + } else { + for _, j := range res { + if p.mem[j.key] == j.plain { + delete(p.mem, j.key) + } + } + } + p.done = upto + p.cond.Broadcast() + p.mu.Unlock() + } +} + +// wait bloque jusqu'à ce que le ticket (et tous ceux d'avant) soit traité. +func (p *convPersister) wait(ticket uint64) { + p.mu.Lock() + for p.done < ticket { + p.cond.Wait() + } + p.mu.Unlock() +} + +// flush attend que tout ce qui a été confié jusqu'ici soit traité. +func (p *convPersister) flush() { + p.mu.Lock() + t := p.queued + p.mu.Unlock() + p.wait(t) +} + +// forget marque une discussion comme supprimée : ses instantanés et résultats en +// attente sont jetés, et plus rien ne la réécrira — sans ça, un instantané en +// vol la ferait renaître (contenu ET entrée d'index) juste après sa suppression. +func (p *convPersister) forget(id string) { + p.mu.Lock() + defer p.mu.Unlock() + p.deleted[id] = true + delete(p.snaps, id) + p.dropToolResLocked(id + ".") +} + +// dropToolRes oublie les résultats en attente d'une discussion. +func (p *convPersister) dropToolRes(prefix string) { + p.mu.Lock() + defer p.mu.Unlock() + p.dropToolResLocked(prefix) +} + +func (p *convPersister) dropToolResLocked(prefix string) { + for k := range p.mem { + if strings.HasPrefix(k, prefix) { + delete(p.mem, k) + } + } + kept := p.res[:0] + for _, j := range p.res { + if !strings.HasPrefix(j.key, prefix) { + kept = append(kept, j) + } + } + p.res = kept +} + +// toolResPending relit un résultat pas encore sur disque. +func (p *convPersister) toolResPending(key string) (string, bool) { + p.mu.Lock() + defer p.mu.Unlock() + s, ok := p.mem[key] + return s, ok +} + +// toolResConv extrait l'identifiant de discussion d'une clé « . ». +func toolResConv(key string) string { + if i := strings.LastIndexByte(key, '.'); i > 0 { + return key[:i] + } + return key +} + +// persistExitWait borne l'attente de l'écrivain avant une sortie : le délai +// d'ouverture de la base (5 s, withDBAt) plus une marge. +const persistExitWait = 6 * time.Second + +// flushPersister attend que les écritures différées soient sur disque, au plus +// d pour ne jamais retenir un arrêt sur une base bloquée. À appeler avant toute +// sortie volontaire du process (redémarrage, « Quitter », exec). +func flushPersister(d time.Duration) { + ok := make(chan struct{}) + go func() { + persistQ.flush() + close(ok) + }() + select { + case <-ok: + case <-time.After(d): + } +} + +// persistCommitByPath regroupe le lot par base (une seule en production) et +// l'écrit. Renvoie les instantanés effectivement écrits. +func persistCommitByPath(snaps []*convSnap, res []toolResJob) ([]*convSnap, error) { + paths := map[string]bool{} + for _, s := range snaps { + paths[s.path] = true + } + for _, j := range res { + paths[j.path] = true + } + var wrote []*convSnap + var firstErr error + for path := range paths { + var ps []*convSnap + var pr []toolResJob + for _, s := range snaps { + if s.path == path { + ps = append(ps, s) + } + } + for _, j := range res { + if j.path == path { + pr = append(pr, j) + } + } + w, err := persistCommit(path, ps, pr) + wrote = append(wrote, w...) + if err != nil && firstErr == nil { + firstErr = err + } + } + return wrote, firstErr +} + +// commitPersistBatch écrit un lot en UNE transaction : les résultats d'outils, +// puis chaque instantané avec la mise à jour de son entrée d'index. Avant, le +// contenu et l'index étaient deux commits (deux fsync), et la relecture de +// l'index hors de tout verrou pouvait écraser un renommage concurrent. +// +// Mémoire chiffrée mais verrouillée : rien n'est écrit en clair, comme +// putStoreBytes — les instantanés sont sautés (le fil reste en RAM), les +// résultats d'outils renvoient une erreur et restent servis depuis la mémoire. +func commitPersistBatch(path string, snaps []*convSnap, res []toolResJob) ([]*convSnap, error) { + // Hors transaction : memEncActive relit la configuration, donc rouvrirait la + // base — bbolt n'est pas rentrant (voir memEncoderNow). + encode, encErr := memEncoderNow() + if len(res) > 0 && encErr != nil { + return nil, encErr + } + now := time.Now().Unix() + convIndexMu.Lock() + defer convIndexMu.Unlock() + cacheBust(bkChat) + cacheBust(bkToolRes) + var wrote []*convSnap + err := withDBAt(path, func(d *bolt.DB) error { + return d.Update(func(tx *bolt.Tx) error { + wrote = wrote[:0] + if len(res) > 0 { + b, err := tx.CreateBucketIfNotExists([]byte(bkToolRes)) + if err != nil { + return err + } + for _, j := range res { + enc, err := encode([]byte(j.plain)) + if err != nil { + return err + } + if err := b.Put([]byte(j.key), enc); err != nil { + return err + } + } + } + if len(snaps) == 0 || encErr != nil { + return nil + } + b, err := tx.CreateBucketIfNotExists([]byte(bkChat)) + if err != nil { + return err + } + // Index illisible (chiffré par une autre clé, abîmé) : on écrit le + // contenu mais on ne touche PAS à l'index — le réécrire depuis une + // lecture vide effacerait toutes les autres discussions de la liste. + var idx []convMeta + idxOK := true + if raw := b.Get([]byte(ckIndex)); len(raw) > 0 { + dec, err := decodeMemContent(raw) + if err != nil || json.Unmarshal(dec, &idx) != nil { + idxOK = false + } + } + for _, s := range snaps { + enc, err := encode(s.body) + if err != nil { + return err + } + if err := b.Put([]byte(convKey(s.id)), enc); err != nil { + return err + } + if idxOK { + idx = convTouchMetaIn(idx, s, now) + } + wrote = append(wrote, s) + } + if !idxOK { + return nil + } + ib, err := json.Marshal(idx) + if err != nil { + return err + } + enc, err := encode(ib) + if err != nil { + return err + } + return b.Put([]byte(ckIndex), enc) + }) + }) + if err != nil { + return nil, err + } + return wrote, nil +} + +// convTouchMetaIn rafraîchit l'entrée d'index d'un instantané : date, nombre +// d'échanges, titre déduit si l'utilisateur n'en a pas choisi. Une discussion +// absente de l'index y est ajoutée — c'est ce qui recolle une discussion créée +// mémoire verrouillée (son entrée n'avait pas pu s'écrire). Une discussion +// SUPPRIMÉE n'arrive jamais ici : l'écrivain l'a écartée (forget). +func convTouchMetaIn(idx []convMeta, s *convSnap, now int64) []convMeta { + for i := range idx { + if idx[i].ID != s.id { + continue + } + idx[i].Updated = now + idx[i].Turns = s.turns + if idx[i].Title == "" { + idx[i].Title = s.title + } + return idx + } + return append(idx, convMeta{ID: s.id, Title: s.title, Created: now, Updated: now, Turns: s.turns, Project: s.project}) +} + +// --- Discussion active, gardée en RAM ----------------------------------------- +// +// convEnsureActive relisait bkChat/active à CHAQUE appel : une ouverture de base +// (3-4 ms sur Windows) pour le mode Code, l'espace de travail, les résultats +// d'outils, la télémétrie… plusieurs fois par tour. Un seul process possède la +// conversation (relay_link.go) et toutes les bascules passent par +// convActivate : le cache ne peut donc pas diverger de la base, exactement +// comme activeProjCache pour le projet actif. +// +// Le cache sert aussi de LIEN entre le contenu en mémoire et son identifiant : +// convActivate le change sous c.mu, en même temps que le contenu (loadFrom). Un +// instantané pris sous c.mu lit donc toujours l'identifiant du contenu qu'il +// sérialise — jamais l'ancien fil sous le nouvel identifiant, même si une +// bascule s'intercale entre deux étapes de persist. + +type convActiveRef struct{ path, id string } + +var ( + // convActiveMu sérialise la création de l'identifiant et les bascules. + // Ordre des verrous : convActiveMu, puis c.mu (loadFrom). Jamais l'inverse — + // d'où la lecture sans verrou du cache chaud. + convActiveMu sync.Mutex + convActiveCache atomic.Pointer[convActiveRef] + // convIndexMu sérialise les lectures-modifications-écritures de l'index : + // celles des opérations de discussion et celle de l'écrivain. + convIndexMu sync.Mutex +) + +// convActiveCached renvoie l'identifiant en cache pour la base courante, ou "". +func convActiveCached() string { + if r := convActiveCache.Load(); r != nil && r.path == dbPath() { + return r.id + } + return "" +} + +// convActiveRemember met l'identifiant en cache. Un pointeur encore chiffré par +// une ancienne version (healPlainStoreKeys) n'est pas mis en cache : il sera +// relu en clair après la réparation. +func convActiveRemember(id string) { + if id == "" || looksEncrypted([]byte(id)) { + convActiveCache.Store(nil) + return + } + convActiveCache.Store(&convActiveRef{path: dbPath(), id: id}) +} + +// convActivate fait de id la discussion active et charge b en mémoire, sous +// convActiveMu : le pointeur en base, le cache et le contenu changent ensemble. +func convActivate(id string, b []byte) { + convActiveMu.Lock() + defer convActiveMu.Unlock() + _ = putStr(bkChat, ckActive, id) + conv.loadFrom(id, b) +} + +// snapshot prend l'instantané à écrire. L'appelant NE doit PAS détenir mu. +func (c *Conversation) snapshot() *convSnap { + // Hors verrou : au premier appel, convEnsureActive lit (voire crée) la clé + // en base. Ensuite le cache est chaud et relu sous c.mu. + id0 := convEnsureActive() + project := activeProjectSlug() + c.mu.Lock() + defer c.mu.Unlock() + id := id0 + if c == conv { + if cid := convActiveCached(); cid != "" { + id = cid + } + } + // 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) + // En plein tour, le journal porte encore un événement par token et chaque + // frappe d'outil : on écrit la forme que produira la fin de tour + // (compactLog), sans toucher au journal en mémoire que suivent les abonnés. + full := c.Log + if c.Generating { + c.Log = compactLog(full) + } + b, err := json.Marshal(c) + c.Log = full + if err != nil { + return nil + } + turns := 0 + for _, ev := range full { + if _, ok := ev.Delta["user"]; ok { + turns++ + } + } + return &convSnap{ + path: dbPath(), + id: id, + seq: persistSeq.Add(1), + body: b, + title: convSummary(c.Messages), + turns: turns, + project: project, + } +} + +// persist enregistre l'état sous sa discussion et ATTEND l'écriture : fin de +// tour, reset, bascules, compactage manuel. Tout ce qui était en attente avant +// lui est écrit aussi. L'appelant NE doit PAS détenir mu. +func (c *Conversation) persist() { + if s := c.snapshot(); s != nil { + persistQ.wait(persistQ.enqueueSnap(s)) + return + } + persistQ.flush() +} + +// persistAsync prend l'instantané tout de suite mais laisse l'écriture à +// l'écrivain : pour les chemins qui précèdent un appel au modèle (StartTurn, +// persistance en cours de tour). L'identifiant de la discussion est résolu ici, +// de façon synchrone — sur une base neuve, generate trouve donc la discussion +// active déjà créée, comme avant. L'appelant NE doit PAS détenir mu. +func (c *Conversation) persistAsync() { + if s := c.snapshot(); s != nil { + persistQ.enqueueSnap(s) + } +} diff --git a/internal/loki/chat_persist_test.go b/internal/loki/chat_persist_test.go new file mode 100644 index 0000000..1696b01 --- /dev/null +++ b/internal/loki/chat_persist_test.go @@ -0,0 +1,425 @@ +package loki + +import ( + "encoding/json" + "fmt" + "reflect" + "strings" + "sync" + "testing" + "time" +) + +// gatePersist remplace l'écriture par une version qu'on retient : chaque lot +// signale son arrivée sur entered puis attend que release soit fermé. Les +// échanges de persistCommit se font écrivain vidé : aucun lot ne le lit alors. +func gatePersist(t *testing.T) (entered chan []*convSnap, release chan struct{}) { + t.Helper() + entered = make(chan []*convSnap, 64) + release = make(chan struct{}) + persistQ.flush() + prev := persistCommit + persistCommit = func(path string, s []*convSnap, r []toolResJob) ([]*convSnap, error) { + entered <- s + <-release + return prev(path, s, r) + } + t.Cleanup(func() { + persistQ.flush() + persistCommit = prev + }) + return entered, release +} + +// waitEntered attend qu'un lot arrive à l'écriture. +func waitEntered(t *testing.T, entered chan []*convSnap) []*convSnap { + t.Helper() + select { + case s := <-entered: + return s + case <-time.After(5 * time.Second): + t.Fatal("aucun lot n'a atteint l'écriture") + return nil + } +} + +// diskMessages relit le nombre de messages et le texte d'une discussion en base. +func diskMessages(t *testing.T, id string) (int, string) { + t.Helper() + b := storeBytes(bkChat, convKey(id)) + if len(b) == 0 { + return -1, "" + } + var st struct { + Messages []Message `json:"messages"` + } + if err := json.Unmarshal(b, &st); err != nil { + t.Fatal(err) + } + return len(st.Messages), string(b) +} + +func addUserMsg(text string) { + conv.mu.Lock() + conv.Messages = append(conv.Messages, Message{Role: "user", Content: text}) + conv.Log = append(conv.Log, LogEvent{Seq: len(conv.Log) + 1, Delta: map[string]any{"user": text}}) + conv.mu.Unlock() +} + +// Des instantanés pris en rafale depuis plusieurs goroutines sont écrits dans +// l'ordre où ils ont été pris : jamais un plus ancien après un plus récent, et +// la base finit sur le dernier état. +func TestPersisterNEcritJamaisUnEtatPlusAncien(t *testing.T) { + testHome(t) + resetConvForTest() + persistQ.flush() + var mu sync.Mutex + var seqs []uint64 + prev := persistCommit + persistCommit = func(path string, s []*convSnap, r []toolResJob) ([]*convSnap, error) { + mu.Lock() + for _, x := range s { + seqs = append(seqs, x.seq) + } + mu.Unlock() + time.Sleep(2 * time.Millisecond) // laisse la file s'accumuler + return prev(path, s, r) + } + t.Cleanup(func() { + persistQ.flush() + persistCommit = prev + }) + + const workers, per = 8, 10 + var wg sync.WaitGroup + for g := 0; g < workers; g++ { + wg.Add(1) + go func(g int) { + defer wg.Done() + for i := 0; i < per; i++ { + addUserMsg(fmt.Sprintf("g%d-%d", g, i)) + conv.persistAsync() + } + }(g) + } + wg.Wait() + persistQ.flush() + + mu.Lock() + defer mu.Unlock() + for i := 1; i < len(seqs); i++ { + if seqs[i] <= seqs[i-1] { + t.Fatalf("instantané %d écrit après %d : un état plus ancien a écrasé un plus récent", seqs[i], seqs[i-1]) + } + } + if n, _ := diskMessages(t, convEnsureActive()); n != workers*per { + t.Fatalf("la base finit sur %d messages, attendu %d (dernier état)", n, workers*per) + } +} + +// persist (fin de tour, reset, bascule) est SYNCHRONE : il attend l'écriture en +// vol, puis la sienne. Au retour de Reset, la base porte déjà le fil vide. +func TestPersistSynchroneAttendLesEcrituresEnVol(t *testing.T) { + testHome(t) + resetConvForTest() + entered, release := gatePersist(t) + addUserMsg("avant le reset") + conv.persistAsync() + waitEntered(t, entered) // l'écrivain est retenu au milieu du lot + + done := make(chan struct{}) + go func() { + conv.Reset() + close(done) + }() + select { + case <-done: + t.Fatal("Reset est revenu avant que l'écriture en vol soit terminée") + case <-time.After(100 * time.Millisecond): + } + close(release) + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("Reset bloqué") + } + if n, _ := diskMessages(t, convEnsureActive()); n != 0 { + t.Fatalf("au retour de Reset, la base doit porter le fil vide (%d message(s))", n) + } +} + +// « Voir plus » marche avant même que le résultat soit sur disque, et la copie +// en mémoire disparaît une fois l'écriture faite. +func TestToolResultServiAvantLEcriture(t *testing.T) { + testHome(t) + entered, release := gatePersist(t) + long := strings.Repeat("sortie\n", 1000) + id := saveToolResult(long) + if id == "" { + t.Fatal("résultat non enregistré") + } + waitEntered(t, entered) + if got, ok := loadToolResult(id); !ok || got != long { + t.Fatal("le résultat doit être servi depuis la mémoire tant qu'il n'est pas écrit") + } + close(release) + persistQ.flush() + if _, ok := persistQ.toolResPending(id); ok { + t.Fatal("la copie en mémoire doit partir une fois le résultat écrit") + } + if got, ok := loadToolResult(id); !ok || got != long { + t.Fatal("résultat perdu après l'écriture") + } +} + +// Un résultat supprimé avec sa discussion avant d'être écrit n'est jamais écrit. +func TestToolResultSupprimeAvantEcriture(t *testing.T) { + testHome(t) + entered, release := gatePersist(t) + saveToolResult(strings.Repeat("a", 3000)) // retient l'écrivain + waitEntered(t, entered) + second := saveToolResult(strings.Repeat("b", 3000)) // en attente derrière + deleteToolResultsFor(toolResConv(second)) + close(release) + persistQ.flush() + if _, ok := loadToolResult(second); ok { + t.Fatal("un résultat supprimé avant son écriture a été écrit quand même") + } +} + +// compactLog rend exactement ce que donne la compaction de fin de tour, sans +// toucher au journal d'origine. +func TestCompactLogPurEgalFinDeTour(t *testing.T) { + log := []LogEvent{ + {Seq: 1, TS: 10, Delta: map[string]any{"user": "q"}}, + {Seq: 2, TS: 11, Delta: map[string]any{"reasoning_content": "je ", "toks": 1}}, + {Seq: 3, TS: 12, Delta: map[string]any{"reasoning_content": "pense", "toks": 1}}, + {Seq: 4, TS: 13, Delta: map[string]any{"tool_used": map[string]any{"name": "bash", "done": false}}}, + {Seq: 5, TS: 14, Delta: map[string]any{"tool_used": map[string]any{"name": "bash", "done": true}}}, + {Seq: 6, TS: 15, Delta: map[string]any{"content": "bon", "toks": 1}}, + {Seq: 7, TS: 16, Delta: map[string]any{"content": "jour", "toks": 1}}, + } + orig := append([]LogEvent(nil), log...) + got := compactLog(log) + if !reflect.DeepEqual(log, orig) { + t.Fatal("compactLog a modifié le journal d'origine") + } + c := newTestConv() + c.Log = append([]LogEvent(nil), log...) + c.compactLogLocked() + if !reflect.DeepEqual(got, c.Log) { + t.Fatalf("compactLog diffère de la fin de tour :\n%v\n%v", got, c.Log) + } + if len(got) != 4 { + t.Fatalf("attendu 4 événements (user, réflexion, outil fini, texte), got %d", len(got)) + } +} + +// En plein tour, la base reçoit le journal compacté ; le journal en mémoire +// (suivi par les abonnés) garde ses événements bruts. +func TestPersistEnPleinTourEcritLeJournalCompacte(t *testing.T) { + testHome(t) + resetConvForTest() + conv.mu.Lock() + conv.Generating = true + for i := 0; i < 50; i++ { + conv.Log = append(conv.Log, LogEvent{Seq: i + 1, Delta: map[string]any{"content": "x", "toks": 1}}) + } + conv.mu.Unlock() + t.Cleanup(func() { + conv.mu.Lock() + conv.Generating = false + conv.mu.Unlock() + }) + conv.persistAsync() + persistQ.flush() + var st struct { + Log []LogEvent `json:"log"` + } + if err := json.Unmarshal(storeBytes(bkChat, convKey(convEnsureActive())), &st); err != nil { + t.Fatal(err) + } + if len(st.Log) != 1 || st.Log[0].Delta["content"] != strings.Repeat("x", 50) { + t.Fatalf("journal écrit non compacté : %d événement(s)", len(st.Log)) + } + conv.mu.Lock() + n := len(conv.Log) + conv.mu.Unlock() + if n != 50 { + t.Fatalf("le journal en mémoire ne doit pas être touché : %d événement(s)", n) + } +} + +// Contenu et entrée d'index partent dans la MÊME écriture. +func TestPersistEcritContenuEtIndexEnUneFois(t *testing.T) { + testHome(t) + resetConvForTest() + persistQ.flush() + var calls int + var mu sync.Mutex + prev := persistCommit + persistCommit = func(path string, s []*convSnap, r []toolResJob) ([]*convSnap, error) { + mu.Lock() + calls++ + mu.Unlock() + return prev(path, s, r) + } + t.Cleanup(func() { + persistQ.flush() + persistCommit = prev + }) + id := convEnsureActive() + addUserMsg("titre de la discussion") + conv.persist() + mu.Lock() + n := calls + mu.Unlock() + if n != 1 { + t.Fatalf("%d écritures pour un persist, attendu 1", n) + } + if c, _ := diskMessages(t, id); c != 1 { + t.Fatalf("contenu non écrit (%d)", c) + } + for _, m := range convIndex() { + if m.ID == id { + if m.Turns != 1 || m.Title != "titre de la discussion" { + t.Fatalf("entrée d'index non rafraîchie : %+v", m) + } + return + } + } + t.Fatal("discussion absente de l'index") +} + +// Supprimer une discussion pendant qu'une écriture la vise ne la fait pas +// renaître : ni contenu, ni entrée d'index. +func TestSuppressionPendantEcritureNeRessusciteRien(t *testing.T) { + testHome(t) + resetConvForTest() + a := convEnsureActive() + addUserMsg("fil A") + conv.persist() + b := convNew() + if err := convSwitch(a); err != nil { + t.Fatal(err) + } + + entered, release := gatePersist(t) + addUserMsg("fil A, suite") + conv.persistAsync() + waitEntered(t, entered) // un lot visant A est en vol + addUserMsg("fil A, encore") + conv.persistAsync() // un autre attend derrière + + done := make(chan error, 1) + go func() { done <- convDelete(a) }() + time.Sleep(50 * time.Millisecond) + close(release) + select { + case err := <-done: + if err != nil { + t.Fatal(err) + } + case <-time.After(5 * time.Second): + t.Fatal("convDelete bloqué") + } + persistQ.flush() + if n, _ := diskMessages(t, a); n != -1 { + t.Fatalf("la discussion supprimée a été réécrite (%d messages)", n) + } + for _, m := range convIndex() { + if m.ID == a { + t.Fatal("la discussion supprimée est revenue dans l'index") + } + } + if got := convEnsureActive(); got != b { + t.Fatalf("active = %q, attendu %q", got, b) + } +} + +// Basculer pendant qu'une écriture est en vol : l'ancien fil est écrit sous SON +// identifiant, jamais sous celui de la discussion qu'on ouvre. +func TestBasculePendantEcritureGardeChaqueFilChezLui(t *testing.T) { + testHome(t) + resetConvForTest() + a := convEnsureActive() + addUserMsg("seulement dans A") + conv.persist() + b := convNew() + addUserMsg("seulement dans B") + conv.persist() + if err := convSwitch(a); err != nil { + t.Fatal(err) + } + + entered, release := gatePersist(t) + addUserMsg("A encore") + conv.persistAsync() + waitEntered(t, entered) + done := make(chan error, 1) + go func() { done <- convSwitch(b) }() + time.Sleep(50 * time.Millisecond) + close(release) + if err := <-done; err != nil { + t.Fatal(err) + } + persistQ.flush() + if _, body := diskMessages(t, b); strings.Contains(body, "dans A") { + t.Fatal("le fil A a été écrit sous l'identifiant de B") + } + if n, body := diskMessages(t, a); n != 2 || !strings.Contains(body, "A encore") { + t.Fatalf("le fil A n'a pas gardé son dernier état (%d messages)", n) + } + conv.mu.Lock() + got := len(conv.Messages) + conv.mu.Unlock() + if got != 1 { + t.Fatalf("la bascule doit charger B (1 message), got %d", got) + } +} + +// Base neuve : l'instantané de StartTurn crée la discussion active TOUT DE +// SUITE, avant l'écriture — generate, qui démarre juste après, voit la même +// discussion (mode Code, critères, espace de travail) que la persistance. +func TestPersistAsyncCreeLaDiscussionActiveAvantGenerate(t *testing.T) { + testHome(t) + resetConvForTest() + entered, release := gatePersist(t) + defer close(release) + addUserMsg("premier message") + conv.persistAsync() + snaps := waitEntered(t, entered) + id := getStr(bkChat, ckActive) + if id == "" { + t.Fatal("la discussion active doit exister dès le retour de persistAsync") + } + if got := convEnsureActive(); got != id { + t.Fatalf("generate verrait %q, la persistance %q", got, id) + } + if len(snaps) != 1 || snaps[0].id != id { + t.Fatalf("l'instantané vise %v, attendu %q", snaps, id) + } +} + +// Base neuve, appels simultanés : un seul identifiant est forgé. +func TestConvEnsureActiveConcurrentForgeUnSeulID(t *testing.T) { + testHome(t) + var wg sync.WaitGroup + ids := make([]string, 8) + for i := range ids { + wg.Add(1) + go func(i int) { + defer wg.Done() + ids[i] = convEnsureActive() + }(i) + } + wg.Wait() + for _, id := range ids[1:] { + if id != ids[0] { + t.Fatalf("identifiants divergents : %v", ids) + } + } + if n := len(convIndex()); n != 1 { + t.Fatalf("%d entrées d'index, attendu 1", n) + } +} diff --git a/internal/loki/chat_sessions.go b/internal/loki/chat_sessions.go index 4e13fc9..a75c597 100644 --- a/internal/loki/chat_sessions.go +++ b/internal/loki/chat_sessions.go @@ -86,6 +86,8 @@ func convIndexForProject(slug string) []convMeta { // tagOrphanConversations rattache au projet donné les discussions qui n'en ont // pas encore (migration, une seule fois : ensuite il n'y a plus d'orphelines). func tagOrphanConversations(slug string) { + convIndexMu.Lock() + defer convIndexMu.Unlock() idx := convIndex() changed := false for i := range idx { @@ -123,8 +125,22 @@ func convIndexSave(idx []convMeta) { // convEnsureActive garantit qu'une discussion active existe, et reprend le fil // unique de l'amont s'il y en a un (migration silencieuse, une seule fois). +// +// L'identifiant est gardé en RAM (chat_persist.go) : seule la première lecture +// ouvre la base. Le test, la lecture et la création se font sous convActiveMu — +// deux appelants simultanés sur une base neuve forgeaient sinon chacun leur +// identifiant, et le second rendait orpheline la discussion du premier. func convEnsureActive() string { + if id := convActiveCached(); id != "" { + return id + } + convActiveMu.Lock() + defer convActiveMu.Unlock() + if id := convActiveCached(); id != "" { + return id + } if id := getStr(bkChat, ckActive); id != "" { + convActiveRemember(id) return id } id := newConvID() @@ -135,7 +151,11 @@ func convEnsureActive() string { } _ = putStr(bkChat, ckActive, id) now := time.Now().Unix() - convIndexSave(append(convIndex(), convMeta{ID: id, Created: now, Updated: now, Project: activeProjectSlug()})) + project := activeProjectSlug() + convIndexMu.Lock() + convIndexSave(append(convIndex(), convMeta{ID: id, Created: now, Updated: now, Project: project})) + convIndexMu.Unlock() + convActiveRemember(id) return id } @@ -143,27 +163,6 @@ func convEnsureActive() string { // dépendance à un générateur aléatoire. func newConvID() string { return fmt.Sprintf("c%d", time.Now().UnixNano()) } -// convTouchMeta rafraîchit les métadonnées de la discussion active après un -// enregistrement : titre (déduit du premier message si l'utilisateur n'en a pas -// choisi), date de modification et nombre d'échanges. -func convTouchMeta(id string, title string, turns int) { - idx := convIndex() - now := time.Now().Unix() - for i := range idx { - if idx[i].ID != id { - continue - } - idx[i].Updated = now - idx[i].Turns = turns - if idx[i].Title == "" { - idx[i].Title = title - } - convIndexSave(idx) - return - } - convIndexSave(append(idx, convMeta{ID: id, Title: title, Created: now, Updated: now, Turns: turns, Project: activeProjectSlug()})) -} - // convSummary lit le premier message utilisateur pour en faire un titre. Sans // message, on laisse vide : l'UI affiche « Nouvelle discussion » et le titre se // posera tout seul au premier échange. @@ -217,8 +216,7 @@ func projectSwitch(slug string) error { return nil } // convIndex rend la plus récemment modifiée en tête. - _ = putStr(bkChat, ckActive, list[0].ID) - conv.loadFrom(storeBytes(bkChat, convKey(list[0].ID))) + convActivate(list[0].ID, storeBytes(bkChat, convKey(list[0].ID))) return nil } @@ -241,12 +239,11 @@ func convSwitch(id string) error { if !found { return fmt.Errorf("discussion introuvable") } - if id == getStr(bkChat, ckActive) { + if id == convEnsureActive() { return nil } conv.persist() // fige la discussion qu'on quitte - _ = putStr(bkChat, ckActive, id) - conv.loadFrom(storeBytes(bkChat, convKey(id))) + convActivate(id, storeBytes(bkChat, convKey(id))) return nil } @@ -256,9 +253,11 @@ func convSwitch(id string) error { func convCreate() string { id := newConvID() now := time.Now().Unix() - convIndexSave(append(convIndex(), convMeta{ID: id, Created: now, Updated: now, Project: activeProjectSlug()})) - _ = putStr(bkChat, ckActive, id) - conv.loadFrom(nil) + project := activeProjectSlug() + convIndexMu.Lock() + convIndexSave(append(convIndex(), convMeta{ID: id, Created: now, Updated: now, Project: project})) + convIndexMu.Unlock() + convActivate(id, nil) return id } @@ -277,6 +276,8 @@ func convRename(id, title string) error { if id == "" || title == "" { return fmt.Errorf("identifiant ou titre manquant") } + convIndexMu.Lock() + defer convIndexMu.Unlock() idx := convIndex() for i := range idx { if idx[i].ID == id { @@ -294,6 +295,7 @@ func convRename(id, title string) error { func convDelete(id string) error { convOpMu.Lock() defer convOpMu.Unlock() + convIndexMu.Lock() idx := convIndex() next := make([]convMeta, 0, len(idx)) found := false @@ -304,10 +306,24 @@ func convDelete(id string) error { } next = append(next, m) } + convIndexMu.Unlock() if !found { return fmt.Errorf("discussion introuvable") } + // Écritures différées (chat_persist.go) : on écarte d'abord tout ce qui + // viserait encore cette discussion, puis on attend l'écriture déjà en vol — + // sinon elle la ferait renaître, contenu et entrée d'index, juste après. + persistQ.forget(id) + persistQ.flush() + convIndexMu.Lock() + next = next[:0] + for _, m := range convIndex() { + if m.ID != id { + next = append(next, m) + } + } convIndexSave(next) + convIndexMu.Unlock() _ = putBytes(bkChat, convKey(id), nil) // Mode code : les jobs d'arrière-plan de la discussion s'arrêtent avec // elle, et ses critères/mode/puce partent avec ses messages. @@ -319,14 +335,13 @@ func convDelete(id string) error { // pour toujours, et plus aucun écran ne permettrait de les retrouver. dropConvFiles(id) deleteToolResultsFor(id) // ses résultats « voir plus » partent avec elle - if id != getStr(bkChat, ckActive) { + if id != convEnsureActive() { return nil } if len(next) == 0 { convCreate() // surtout pas convNew : il réenregistrerait la supprimée return nil } - _ = putStr(bkChat, ckActive, next[0].ID) - conv.loadFrom(storeBytes(bkChat, convKey(next[0].ID))) + convActivate(next[0].ID, storeBytes(bkChat, convKey(next[0].ID))) return nil } diff --git a/internal/loki/code_verify.go b/internal/loki/code_verify.go index f73eb42..2b5002c 100644 --- a/internal/loki/code_verify.go +++ b/internal/loki/code_verify.go @@ -195,7 +195,9 @@ func (c *Conversation) runBuilderTurn(ctx context.Context, caps Caps, temperatur } } c.mu.Unlock() - c.persist() + // En plein tour : l'étape suivante part vers le modèle juste après, le fsync + // ne doit pas l'attendre (chat_persist.go). + c.persistAsync() return asked } diff --git a/internal/loki/mem_store.go b/internal/loki/mem_store.go index 22f1448..37bc302 100644 --- a/internal/loki/mem_store.go +++ b/internal/loki/mem_store.go @@ -174,6 +174,9 @@ var encryptedBuckets = []string{bkChat, bkRecall, bkTracker, bkToolRes} // reencryptChatStores (re)chiffre les buckets de conversation. Exige la DEK. func reencryptChatStores() error { + // Écritures différées (chat_persist.go) d'abord : écrites pendant le + // passage, elles pourraient être écrasées par la relecture d'avant. + persistQ.flush() healPlainStoreKeys() for _, b := range encryptedBuckets { if err := reencryptBucket(b); err != nil { @@ -185,6 +188,7 @@ func reencryptChatStores() error { // decryptChatStores remet en clair les buckets de conversation. Exige la DEK. func decryptChatStores() error { + persistQ.flush() // même raison que reencryptChatStores for _, b := range encryptedBuckets { if err := decryptBucket(b); err != nil { return err diff --git a/internal/loki/store.go b/internal/loki/store.go index 46369bd..607319e 100644 --- a/internal/loki/store.go +++ b/internal/loki/store.go @@ -62,8 +62,14 @@ func dbPath() string { return filepath.Join(LokiHome(), "loki.db") } // withDB ouvre la base, exécute fn, puis referme — toujours, même en erreur. func withDB(fn func(*bolt.DB) error) error { - path := dbPath() + return withDBAt(dbPath(), fn) +} +// withDBAt : withDB sur une base désignée. Sert à l'écrivain de la discussion +// (chat_persist.go), qui écrit APRÈS coup et doit viser la base de l'instant où +// l'instantané a été pris — pas celle que désignerait LOKI_HOME au moment de +// l'écriture (les tests en changent d'un test à l'autre). +func withDBAt(path string, fn func(*bolt.DB) error) error { dbMu.Lock() defer dbMu.Unlock() if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil { diff --git a/internal/loki/store_test.go b/internal/loki/store_test.go index 205b120..b3582d1 100644 --- a/internal/loki/store_test.go +++ b/internal/loki/store_test.go @@ -11,6 +11,9 @@ func testHome(t *testing.T) string { t.Setenv("LOKI_HOME", home) firewallInert = true t.Cleanup(func() { firewallInert = false }) + // Écritures différées (chat_persist.go) : vidées AVANT que TempDir ne + // supprime la base — les cleanups passent dans l'ordre inverse. + t.Cleanup(persistQ.flush) return home } diff --git a/internal/loki/sys_platform_unix.go b/internal/loki/sys_platform_unix.go index 68b5f74..5ffecee 100644 --- a/internal/loki/sys_platform_unix.go +++ b/internal/loki/sys_platform_unix.go @@ -254,6 +254,9 @@ func refreshToolPath() {} // supervises llama-server directly, as the old start.sh did with `exec`). // args[0] must be the binary path. func execServer(bin string, args []string) error { + // exec remplace le process sans rien exécuter après : ce qui attendait + // l'écrivain de la discussion (chat_persist.go) serait perdu. + flushPersister(persistExitWait) return syscall.Exec(bin, args, os.Environ()) } diff --git a/internal/loki/sys_restart_windows.go b/internal/loki/sys_restart_windows.go index 6fa735b..52ae461 100644 --- a/internal/loki/sys_restart_windows.go +++ b/internal/loki/sys_restart_windows.go @@ -42,6 +42,7 @@ func scheduleAppRestart() (bool, string) { go func() { time.Sleep(1500 * time.Millisecond) // laisser la réponse atteindre le navigateur + flushPersister(persistExitWait) // écritures différées de la discussion os.Exit(0) }() return true, "Loki redémarre — la page se reconnectera toute seule dans quelques secondes." diff --git a/internal/loki/sys_tray_darwin.go b/internal/loki/sys_tray_darwin.go index 2b69504..ffa9ae3 100644 --- a/internal/loki/sys_tray_darwin.go +++ b/internal/loki/sys_tray_darwin.go @@ -50,6 +50,7 @@ func runTray(url string) { // que seul root peut piloter. Fermer une interface n'a pas à arrêter un // service système. Sous Windows, où l'app est propriétaire du moteur, // « Quitter » l'arrête bel et bien (voir sys_tray_windows.go). + flushPersister(persistExitWait) // écritures différées de la discussion os.Exit(0) }) } diff --git a/internal/loki/sys_tray_windows.go b/internal/loki/sys_tray_windows.go index eae5fc6..53133c6 100644 --- a/internal/loki/sys_tray_windows.go +++ b/internal/loki/sys_tray_windows.go @@ -53,6 +53,7 @@ func runTray(url string) { // Linux et macOS c'est systemd ou launchd, et fermer une interface n'a // pas à arrêter un service système. _ = svcStop(false) + flushPersister(persistExitWait) // écritures différées de la discussion os.Exit(0) }) } diff --git a/internal/loki/tool_results.go b/internal/loki/tool_results.go index 25376eb..847d560 100644 --- a/internal/loki/tool_results.go +++ b/internal/loki/tool_results.go @@ -45,13 +45,23 @@ func toolResID(sid string) string { // saveToolResult enregistre un résultat complet de la discussion ACTIVE et // renvoie son id ("" en cas d'échec : l'appelant envoie alors le résultat // entier dans le flux). +// +// L'écriture est confiée à l'écrivain de la discussion (chat_persist.go) : un +// tour d'agent qui coupe vingt résultats payait vingt fsync entre l'outil et +// l'étape suivante. Le résultat est servi depuis la mémoire jusqu'à ce qu'il +// soit sur disque (loadToolResult). func saveToolResult(result string) string { - id := toolResID(getStr(bkChat, ckActive)) - // Chiffrement actif mais mémoire verrouillée : putStoreBytes refuse d'écrire - // en clair, et l'appelant envoie alors le résultat entier dans le flux. - if id == "" || putStoreBytes(bkToolRes, id, []byte(result)) != nil { + id := toolResID(convEnsureActive()) + // Chiffrement actif mais mémoire verrouillée : on refuse d'écrire en clair, + // et l'appelant envoie alors le résultat entier dans le flux. Vérifié ICI, + // pas à l'écriture : un id rendu doit désigner un résultat qu'on écrira. + if id == "" { return "" } + if _, err := memEncoderNow(); err != nil { + return "" + } + persistQ.enqueueToolRes(toolResJob{path: dbPath(), key: id, plain: result}) if toolResWrites.Add(1)%toolResPruneGap == 0 { go pruneToolResults() } @@ -63,6 +73,9 @@ func loadToolResult(id string) (string, bool) { if strings.ContainsAny(id, "/\\") { return "", false } + if s, ok := persistQ.toolResPending(id); ok { + return s, true + } b, err := getStoreBytesErr(bkToolRes, id) if err != nil || b == nil { return "", false @@ -76,6 +89,7 @@ func deleteToolResultsFor(sid string) { return } prefix := sid + "." + persistQ.dropToolRes(prefix) _ = update(bkToolRes, func(b *bolt.Bucket) error { c := b.Cursor() var keys [][]byte diff --git a/internal/loki/tool_results_test.go b/internal/loki/tool_results_test.go index 5966bf5..0c639a5 100644 --- a/internal/loki/tool_results_test.go +++ b/internal/loki/tool_results_test.go @@ -55,6 +55,7 @@ func TestToolResultsChiffres(t *testing.T) { if id == "" { t.Fatal("résultat non enregistré") } + persistQ.flush() // l'écriture est différée (chat_persist.go) if !looksEncrypted(getBytes(bkToolRes, id)) { t.Fatal("résultat en clair alors que la mémoire est chiffrée") }