mirror of
https://github.com/R0m1k3/Loki.git
synced 2026-10-11 17:26:57 +02:00
Persistance : le fsync de la discussion ne retarde plus le premier token
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 <noreply@anthropic.com>
This commit is contained in:
1 parent
9b2e974942
commit
c3005f5b56
14 files changed
+1082
-80
No files matched your search
@@ -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).
|
||||
|
||||
@@ -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é « <id>.<aléa> ».
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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())
|
||||
}
|
||||
|
||||
|
||||
@@ -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."
|
||||
|
||||
@@ -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)
|
||||
})
|
||||
}
|
||||
@@ -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)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user