mirror of
https://github.com/R0m1k3/Loki.git
synced 2026-10-11 17:26:57 +02:00
Démarrage, captures, postes distants : correctifs portés d'AJEAN
- Au démarrage, une lecture ratée de la base (verrou bbolt transitoire pendant le chevauchement des process au redémarrage du conteneur) était confondue avec une conversation absente. Pire que l'amont chez nous : getStr rend "" dans les deux cas, donc convEnsureActive forgeait un NOUVEL identifiant et l'écrasait — le fil en cours devenait orphelin, en silence. On sonde désormais la base avec son erreur AVANT toute écriture, on réessaie quatre fois, et un échec durable est journalisé sans que rien ne soit touché. - Une capture d'écran disparaissait sans laisser de trace : stripImageParts aplatissait le message en ne gardant que sa légende. Le marqueur imageLostMarker rend la perte VISIBLE, pour que le modèle sache reprendre une capture au lieu de la redécrire de mémoire. - maxLogEvents 20000 → 200000 : un seul tour à très long raisonnement tronquait déjà le journal de rejeu, et l'utilisateur perdait le début de sa conversation à l'écran. - Keepalive WebSocket des postes distants (ping toutes les 25 s, des DEUX côtés). Sans lui, un poste au repos était coupé au bout de ~60-100 s par les intermédiaires qui ferment les canaux inactifs, puis reconnecté après backoff — les déconnexions à répétition sur tous les postes. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
1 parent
2d4f5378c0
commit
c45ba615fa
6 files changed
+308
-11
No files matched your search
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
@@ -28,8 +29,12 @@ import (
|
||||
// qu'une poignée d'événements au lieu d'un par token → 20000 couvre des
|
||||
// centaines de tours. La marge sert surtout à absorber UN tour en cours (streamé
|
||||
// token par token) avant sa coalescence : un très gros tour (long raisonnement +
|
||||
// réponse) ne doit pas se faire tronquer le début avant d'être compacté.
|
||||
const maxLogEvents = 20000
|
||||
// réponse) ne doit pas se faire tronquer le début avant d'être compacté. 20000
|
||||
// n'y suffisait pas — un seul tour à très long raisonnement tronquait déjà le
|
||||
// journal de rejeu, et l'utilisateur perdait le DÉBUT de sa conversation à
|
||||
// l'écran. Le journal ne pèse que de la mémoire d'affichage, jamais du contexte
|
||||
// modèle : la marge est bon marché.
|
||||
const maxLogEvents = 200000
|
||||
|
||||
// LogEvent = un événement d'affichage rejouable (un delta SSE + son numéro de
|
||||
// séquence monotone + un horodatage serveur en ms). Le TS permet au client de
|
||||
@@ -74,12 +79,49 @@ var conv = func() *Conversation {
|
||||
// modèle. En clair : cette machine déchiffre déjà pour lancer le modèle, le
|
||||
// relais reste aveugle — la persister ici ne change rien à la posture E2E.
|
||||
|
||||
// loadConvAttempts / loadConvRetryWait : au redémarrage du conteneur, l'ancien
|
||||
// process peut encore tenir le verrou bbolt quelques centaines de ms pendant que
|
||||
// le nouveau démarre. Quatre essais espacés de 250 ms couvrent ce chevauchement
|
||||
// sans retarder perceptiblement un démarrage normal (où le premier essai passe).
|
||||
const (
|
||||
loadConvAttempts = 4
|
||||
loadConvRetryWait = 250 * time.Millisecond
|
||||
)
|
||||
|
||||
// LoadConversation recharge l'état persisté au démarrage du process. Sans état
|
||||
// enregistré (première fois) on part d'une conversation vide.
|
||||
//
|
||||
// Une lecture RATÉE n'est PAS une absence, et les confondre coûte cher ici :
|
||||
// getStr rend "" aussi bien pour « clé absente » que pour « base inaccessible »,
|
||||
// donc sur un verrou transitoire convEnsureActive croyait la discussion active
|
||||
// inexistante, forgeait un NOUVEL identifiant et l'écrasait — le fil en cours
|
||||
// devenait orphelin, en silence. On sonde donc la base avec son erreur AVANT
|
||||
// toute écriture, on réessaie, et en cas d'échec durable on le dit et on ne
|
||||
// touche à rien.
|
||||
func LoadConversation() {
|
||||
// convEnsureActive reprend au passage le fil unique des versions
|
||||
// précédentes (clé « conversation ») comme première discussion.
|
||||
b := getBytes(bkChat, convKey(convEnsureActive()))
|
||||
var b []byte
|
||||
var err error
|
||||
for attempt := 0; attempt < loadConvAttempts; attempt++ {
|
||||
// Sonde : lire la clé de la discussion active en gardant l'erreur. Tant
|
||||
// qu'elle échoue, convEnsureActive ne doit surtout pas être appelée.
|
||||
if _, err = getBytesErr(bkChat, ckActive); err != nil {
|
||||
time.Sleep(loadConvRetryWait)
|
||||
continue
|
||||
}
|
||||
// convEnsureActive reprend au passage le fil unique des versions
|
||||
// précédentes (clé « conversation ») comme première discussion.
|
||||
if b, err = getBytesErr(bkChat, convKey(convEnsureActive())); err == nil {
|
||||
break
|
||||
}
|
||||
time.Sleep(loadConvRetryWait)
|
||||
}
|
||||
if err != nil {
|
||||
// Toujours en échec : on le DIT au lieu de repartir à vide en silence, et
|
||||
// on abandonne le chargement sans rien écrire — un serveur qui refuse de
|
||||
// démarrer serait pire, et l'historique sur disque reste intact.
|
||||
fmt.Fprintf(os.Stderr, "[conv] base illisible au démarrage (%v) — aucune discussion chargée ; rien n'a été écrasé, l'historique est intact\n", err)
|
||||
return
|
||||
}
|
||||
if len(b) == 0 {
|
||||
return
|
||||
}
|
||||
@@ -479,9 +521,12 @@ func (c *Conversation) generate(ctx context.Context, caps Caps, temperature floa
|
||||
// Prompt système personnalisé (UI → /api/sysprompt, fichier côté serveur).
|
||||
// Injecté seulement dans la vue envoyée au modèle, jamais persisté dans
|
||||
// c.Messages : modifiable à chaud, effet dès le tour suivant.
|
||||
final := msgs
|
||||
// Contexte du projet actif (description, index mémoire, index des trackers),
|
||||
// même traitement : injecté dans la vue envoyée, jamais persisté — voir
|
||||
// projectSystemMessages.
|
||||
final := append(projectSystemMessages(), msgs...)
|
||||
if sp := readSysPrompt(); sp != "" {
|
||||
final = append([]Message{{Role: "system", Content: sp}}, msgs...)
|
||||
final = append([]Message{{Role: "system", Content: sp}}, final...)
|
||||
}
|
||||
|
||||
// newBase : vue modèle publiée par une compaction survenue PENDANT le tour.
|
||||
|
||||
@@ -0,0 +1,146 @@
|
||||
package loki
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
bolt "go.etcd.io/bbolt"
|
||||
)
|
||||
|
||||
// resetConvForTest vide l'état en RAM pour forcer un VRAI rechargement depuis la
|
||||
// base, comme au démarrage d'un process qui ne connaît encore rien. `conv` est un
|
||||
// global du paquet : sans ça, un test verrait l'état laissé par le précédent.
|
||||
func resetConvForTest() {
|
||||
conv.mu.Lock()
|
||||
conv.Messages, conv.Log, conv.Seq, conv.CtxUsed = nil, nil, 0, 0
|
||||
conv.mu.Unlock()
|
||||
}
|
||||
|
||||
// TestLoadConversationReessayeSousContention : la panne d'origine est un verrou
|
||||
// bbolt transitoire au redémarrage du conteneur (l'ancien process tient encore
|
||||
// le fichier pendant que le nouveau démarre), confondu avec « rien à charger ».
|
||||
//
|
||||
// Réserve honnête : ce test ne reproduit PAS la panne à coup sûr. withDB ouvre
|
||||
// déjà la base avec Options.Timeout de 5 s, qui absorbe seul une contention de
|
||||
// 400 ms — ce test passe donc aussi sur le code d'AVANT le correctif. Il est
|
||||
// gardé comme filet de non-régression sur le nouveau chemin (getBytesErr + boucle
|
||||
// de tentatives) : ce chemin ne doit ni perdre la conversation, ni se bloquer.
|
||||
func TestLoadConversationReessayeSousContention(t *testing.T) {
|
||||
t.Setenv("LOKI_HOME", t.TempDir())
|
||||
|
||||
// 1. Écrit une conversation par l'API normale (elle crée aussi la clé active).
|
||||
resetConvForTest()
|
||||
conv.mu.Lock()
|
||||
conv.Messages = []Message{{Role: "user", Content: "bonjour"}}
|
||||
conv.mu.Unlock()
|
||||
conv.persist()
|
||||
activeID := getStr(bkChat, ckActive)
|
||||
if activeID == "" {
|
||||
t.Fatal("persist n'a pas créé de discussion active — le montage du test est faux")
|
||||
}
|
||||
|
||||
// 2. État en RAM vidé : le rechargement doit tout reprendre depuis la base.
|
||||
resetConvForTest()
|
||||
|
||||
// 3. Verrou exclusif tenu une fenêtre courte dans une goroutine séparée : la
|
||||
// base est momentanément indisponible, exactement comme au chevauchement.
|
||||
db, err := bolt.Open(dbPath(), 0o600, &bolt.Options{Timeout: 2 * time.Second})
|
||||
if err != nil {
|
||||
t.Fatalf("impossible de prendre le verrou pour le test : %v", err)
|
||||
}
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
time.Sleep(400 * time.Millisecond)
|
||||
db.Close()
|
||||
}()
|
||||
|
||||
LoadConversation()
|
||||
wg.Wait()
|
||||
|
||||
conv.mu.Lock()
|
||||
got := len(conv.Messages)
|
||||
conv.mu.Unlock()
|
||||
if got != 1 {
|
||||
t.Fatalf("conversation perdue malgré une contention temporaire : %d message(s), attendu 1", got)
|
||||
}
|
||||
if id := getStr(bkChat, ckActive); id != activeID {
|
||||
t.Fatalf("la discussion active a changé (%q → %q) : le fil d'origine est orphelin", activeID, id)
|
||||
}
|
||||
}
|
||||
|
||||
// TestLoadConversationBaseIllisibleNeForgePasDeDiscussion garde l'invariant qui
|
||||
// justifie le correctif : convEnsureActive ÉCRIT (elle forge un identifiant et
|
||||
// l'enregistre sous bkChat/active). Elle ne doit donc jamais être appelée sur la
|
||||
// foi d'une lecture ratée — sinon un verrou durable suffisait à mettre le fil en
|
||||
// cours de côté au profit d'une discussion neuve, en silence.
|
||||
//
|
||||
// Réserve honnête, là encore : avec une base illisible l'écriture échoue elle
|
||||
// aussi, donc le disque finit identique dans les deux versions. Ce que ce test
|
||||
// verrouille, c'est le CONTRAT — sortir sans rien écrire, et sans faire attendre
|
||||
// indéfiniment — pas la reproduction du dégât.
|
||||
func TestLoadConversationBaseIllisibleNeForgePasDeDiscussion(t *testing.T) {
|
||||
home := t.TempDir()
|
||||
t.Setenv("LOKI_HOME", home)
|
||||
|
||||
// Fichier de base invalide : bolt.Open échoue tout de suite (pas d'attente de
|
||||
// verrou), ce qui exerce le chemin d'erreur sans allonger le test.
|
||||
path := filepath.Join(home, "loki.db")
|
||||
if err := os.WriteFile(path, []byte("ceci n'est pas une base bbolt"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
before, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
resetConvForTest()
|
||||
start := time.Now()
|
||||
LoadConversation()
|
||||
elapsed := time.Since(start)
|
||||
|
||||
// Budget : loadConvAttempts tentatives espacées de loadConvRetryWait, plus une
|
||||
// marge. Au-delà, la boucle de reprise retarderait le démarrage du serveur.
|
||||
if max := time.Duration(loadConvAttempts)*loadConvRetryWait + 2*time.Second; elapsed > max {
|
||||
t.Fatalf("démarrage retardé de %v par une base illisible (budget %v)", elapsed, max)
|
||||
}
|
||||
after, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if string(before) != string(after) {
|
||||
t.Fatal("la base a été modifiée alors qu'elle n'a même pas pu être lue")
|
||||
}
|
||||
conv.mu.Lock()
|
||||
n := len(conv.Messages)
|
||||
conv.mu.Unlock()
|
||||
if n != 0 {
|
||||
t.Fatalf("%d message(s) chargés depuis une base illisible", n)
|
||||
}
|
||||
}
|
||||
|
||||
// TestLoadConversationAbsenceLegitimeResteRapide : le correctif ne doit pas avoir
|
||||
// transformé le cas NORMAL (base saine, rien d'enregistré encore — premier
|
||||
// démarrage) en attente inutile. Une seule tentative doit suffire.
|
||||
func TestLoadConversationAbsenceLegitimeResteRapide(t *testing.T) {
|
||||
t.Setenv("LOKI_HOME", t.TempDir())
|
||||
resetConvForTest()
|
||||
|
||||
start := time.Now()
|
||||
LoadConversation()
|
||||
elapsed := time.Since(start)
|
||||
|
||||
if elapsed >= loadConvRetryWait {
|
||||
t.Fatalf("absence légitime traitée comme une erreur : %v d'attente, aucune reprise ne devait être déclenchée", elapsed)
|
||||
}
|
||||
conv.mu.Lock()
|
||||
n := len(conv.Messages)
|
||||
conv.mu.Unlock()
|
||||
if n != 0 {
|
||||
t.Fatalf("%d message(s) chargés alors que rien n'avait été enregistré", n)
|
||||
}
|
||||
}
|
||||
@@ -154,6 +154,11 @@ func engineSeesImages() bool {
|
||||
// où un base64 de plusieurs dizaines de milliers de tokens était rejoué à
|
||||
// chaque tour jusqu'à dépasser définitivement le contexte ;
|
||||
// - garantir l'invariant à l'avenir, quel que soit le chemin d'écriture.
|
||||
//
|
||||
// Le texte restant reçoit imageLostMarker : sans lui, une conversation rouverte
|
||||
// ne garde que la légende de la capture (« Capture … ») et le modèle n'a plus
|
||||
// aucun indice qu'il A VU une image — il la redécrit de mémoire au lieu d'en
|
||||
// reprendre une. Même raison que dans msgText, autre bout de la chaîne.
|
||||
func stripImageParts(msgs []Message) []Message {
|
||||
for i, m := range msgs {
|
||||
parts, ok := m.Content.([]any)
|
||||
@@ -176,7 +181,7 @@ func stripImageParts(msgs []Message) []Message {
|
||||
}
|
||||
}
|
||||
if dropped {
|
||||
msgs[i].Content = strings.Join(texts, "\n")
|
||||
msgs[i].Content = strings.Join(texts, "\n") + imageLostMarker
|
||||
}
|
||||
}
|
||||
return msgs
|
||||
|
||||
@@ -139,7 +139,9 @@ func TestScreenshotSuitLEtatDeLaVision(t *testing.T) {
|
||||
// Un base64 d'image persisté dans l'historique est rejoué à chaque tour : vu en
|
||||
// production, 55 000 tokens de requête pour 32 768 de contexte — plus aucun
|
||||
// tour ne passait. stripImageParts guérit les conversations existantes en
|
||||
// retirant les parties image et en aplatissant le texte restant.
|
||||
// retirant les parties image et en aplatissant le texte restant — en laissant
|
||||
// imageLostMarker, pour que le modèle sache qu'une image A ÉTÉ montrée et qu'il
|
||||
// doit en reprendre une plutôt que la redécrire de mémoire.
|
||||
func TestStripImagePartsGueritLHistorique(t *testing.T) {
|
||||
msgs := []Message{
|
||||
{Role: "user", Content: "bonjour"}, // simple chaîne : intouchée
|
||||
@@ -153,8 +155,40 @@ func TestStripImagePartsGueritLHistorique(t *testing.T) {
|
||||
t.Fatalf("message texte modifié : %#v", out[0].Content)
|
||||
}
|
||||
got, ok := out[1].Content.(string)
|
||||
if !ok || got != "Voici la capture demandée." {
|
||||
t.Fatalf("le message multimodal doit être aplati en texte sans l'image, obtenu %#v", out[1].Content)
|
||||
if !ok || got != "Voici la capture demandée."+imageLostMarker {
|
||||
t.Fatalf("le message multimodal doit être aplati en texte, sans l'image mais avec le marqueur de perte, obtenu %#v", out[1].Content)
|
||||
}
|
||||
if strings.Contains(got, "base64") {
|
||||
t.Fatalf("le base64 de l'image ne doit plus apparaître : %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
// msgText est le point de passage de TOUTE la chaîne de compaction (estimation
|
||||
// du contexte, transcript donné au résumeur). S'il n'extrait que les parties
|
||||
// `text`, un message qui porte une capture ne pèse que sa légende et le résumeur
|
||||
// ignore jusqu'à l'existence de l'image : elle disparaît sans laisser de trace.
|
||||
func TestMsgTextSignaleLaPresenceDUneImage(t *testing.T) {
|
||||
// Message VIVANT tel que construit par screenshotImageMessage ([]map[string]any).
|
||||
live := Message{Role: "user", Content: []map[string]any{
|
||||
{"type": "text", "text": "Capture demandée :"},
|
||||
{"type": "image_url", "image_url": map[string]any{"url": "data:image/jpeg;base64,AAAA"}},
|
||||
}}
|
||||
// Même message RELU depuis le JSON persisté ([]any de map génériques).
|
||||
reloaded := Message{Role: "user", Content: []any{
|
||||
map[string]any{"type": "text", "text": "Capture demandée :"},
|
||||
map[string]any{"type": "image_url", "image_url": map[string]any{"url": "data:image/jpeg;base64,AAAA"}},
|
||||
}}
|
||||
for name, m := range map[string]Message{"vivant": live, "relu": reloaded} {
|
||||
got := msgText(m)
|
||||
if !strings.Contains(got, "Capture demandée :") {
|
||||
t.Fatalf("%s : le texte de la légende doit être conservé, obtenu %q", name, got)
|
||||
}
|
||||
if !strings.Contains(got, imageLostMarker) {
|
||||
t.Fatalf("%s : la présence de l'image doit être signalée, obtenu %q", name, got)
|
||||
}
|
||||
if strings.Contains(got, "base64") {
|
||||
t.Fatalf("%s : le base64 ne doit jamais entrer dans le transcript : %q", name, got)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -25,6 +25,12 @@ const (
|
||||
nodePairCodeTTL = 10 * time.Minute
|
||||
nodeCallTimeout = 5 * time.Minute // un shell distant peut être long
|
||||
nodeHelloWait = 15 * time.Second
|
||||
// nodePingInterval : cadence du keepalive WebSocket poste↔serveur. 25 s tient
|
||||
// sous les délais d'inactivité usuels des proxys et relais (~60-100 s).
|
||||
nodePingInterval = 25 * time.Second
|
||||
// nodePingTimeout : au-delà, le poste ne répond plus — on ferme au lieu de
|
||||
// laisser une connexion morte occuper le registre.
|
||||
nodePingTimeout = 10 * time.Second
|
||||
)
|
||||
|
||||
// pairedNode est un poste appairé, tel que persisté dans bkState["nodes"].
|
||||
@@ -279,6 +285,34 @@ func handleNodeWS(w http.ResponseWriter, r *http.Request) {
|
||||
defer nodeUnregister(nc)
|
||||
touchNodeSeen(pn.ID)
|
||||
|
||||
// Keepalive : un ping WebSocket périodique garde la connexion vivante à travers
|
||||
// les intermédiaires qui coupent les canaux INACTIFS (relais, Caddy, Cloudflare :
|
||||
// ~60-100 s sans trafic). Sans lui, un poste au repos — l'IA ne le sollicite pas
|
||||
// en permanence — était déconnecté en boucle puis reconnecté après backoff.
|
||||
// coder/websocket sérialise Ping avec les écritures de données en interne : pas
|
||||
// de course avec nc.send. Un ping sans réponse ferme la connexion, ce qui
|
||||
// débloque le Read de la boucle ci-dessous.
|
||||
pingCtx, stopPing := context.WithCancel(ctx)
|
||||
defer stopPing()
|
||||
go func() {
|
||||
t := time.NewTicker(nodePingInterval)
|
||||
defer t.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-pingCtx.Done():
|
||||
return
|
||||
case <-t.C:
|
||||
pc, cancel := context.WithTimeout(pingCtx, nodePingTimeout)
|
||||
err := c.Ping(pc)
|
||||
cancel()
|
||||
if err != nil {
|
||||
_ = c.Close(websocket.StatusGoingAway, "ping sans réponse")
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
// Boucle de lecture : des « result » chiffrés qui débloquent l'appel en attente.
|
||||
for {
|
||||
m, err := nc.readMsg(ctx)
|
||||
|
||||
@@ -36,6 +36,13 @@ const (
|
||||
maxTimeoutSec = 300
|
||||
maxOutput = 8000 // caractères de stdout/stderr renvoyés
|
||||
readMax = 100000 // caractères max d'une lecture de fichier
|
||||
// nodePingInterval : cadence du keepalive WebSocket. 25 s tient sous les délais
|
||||
// d'inactivité usuels des proxys et relais (~60-100 s). Doit rester alignée sur
|
||||
// la constante homonyme côté serveur (node_server.go).
|
||||
nodePingInterval = 25 * time.Second
|
||||
// nodePingTimeout : au-delà, le serveur ne répond plus — on ferme, et Run
|
||||
// reconnecte proprement au lieu de garder une connexion morte.
|
||||
nodePingTimeout = 10 * time.Second
|
||||
)
|
||||
|
||||
// Config persiste l'appairage et les préférences locales du poste.
|
||||
@@ -245,6 +252,32 @@ func session(ctx context.Context, cfg Config, quiet bool) error {
|
||||
if !quiet {
|
||||
fmt.Printf("[node] connecté ✓ (canal chiffré, en attente des demandes de l'IA)\n")
|
||||
}
|
||||
// Keepalive : un ping WebSocket périodique empêche les intermédiaires (relais,
|
||||
// Caddy, Cloudflare) de couper la connexion quand l'IA ne sollicite pas le
|
||||
// poste. Sans lui, un canal inactif tombait au bout de ~60-100 s, la session se
|
||||
// terminait et Run reconnectait après backoff — d'où les déconnexions à
|
||||
// répétition sur tous les postes. coder/websocket sérialise Ping avec les
|
||||
// écritures de données : pas de course avec sendEnc.
|
||||
pingCtx, stopPing := context.WithCancel(ctx)
|
||||
defer stopPing()
|
||||
go func() {
|
||||
t := time.NewTicker(nodePingInterval)
|
||||
defer t.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-pingCtx.Done():
|
||||
return
|
||||
case <-t.C:
|
||||
pc, cancel := context.WithTimeout(pingCtx, nodePingTimeout)
|
||||
err := c.Ping(pc)
|
||||
cancel()
|
||||
if err != nil {
|
||||
_ = c.Close(websocket.StatusGoingAway, "ping sans réponse")
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
for {
|
||||
var fr nodewire.Frame
|
||||
if err := wsjson.Read(ctx, c, &fr); err != nil {
|
||||
|
||||
Reference in new issue
Block a user