Télémétrie : chaque complétion dit enfin ce qu'elle a repris du cache, recalculé et accepté du brouillon

Jusqu'ici, impossible de savoir si une ligne volatile s'était glissée dans
le prompt, si une mise à jour du moteur avait cassé les points de reprise
ou si l'acceptation du MTP s'était effondrée : le chunk final de llama.cpp
porte ces chiffres (cache_n, draft_n, cached_tokens), on les jetait.

Lecture seule : aucun champ ajouté aux requêtes, et le comptage du
contexte (usage.prompt_tokens → CtxUsed → compaction) ne change pas.

- perf_log.go : décodage à part et tolérant (perfWire) — un compteur mal
  typé ne jette plus le chunk final ni ne fait échouer une compaction ;
  cached_tokens, strictement entier dans streamChunk, y déménage aussi.
- Inconnu n'est pas 0 : pointeurs, et l'UI garde « prefill N tok » quand
  le moteur ne dit rien de son cache.
- Nature (main, subagent, verify, task, compact, bench, foreign) portée
  par le contexte, pas par Caps ; TTFT au premier delta de tout type.
- Perte de cache = total précédent − cache_n, seulement quand les messages
  prolongent strictement ceux de la complétion précédente de la même
  discussion et de la même nature (empreintes cumulées) ; jamais sur un
  flux coupé. Ce qui s'est intercalé (sous-agent, client /v1, bench) est
  nommé.
- Anneau de 5000 entrées en mémoire, sans texte ; GET /api/perf/summary
  (derrière la clé) : taux de cache, recalcul par tour et par nature,
  médianes pp/tg par profondeur, acceptation du brouillon.
- Ligne [perf] sur stderr seulement avec LOKI_PERF_LOG.
- UI : « prefill X nouveaux / Y en cache », ambre au-delà d'un seuil
  relevé sur un modèle hybride (espacement des points de reprise).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
MichaelandClaude Opus 5.5 committed 2026-10-04 00:29:17 +02:00
1 parent 3ff56e01db
commit 80b7793773
14 files changed
+1032 -25

No files matched your search

+10 -2
View File
@@ -676,10 +676,18 @@ Write the summary in the SAME language as the conversation.`
}
return "", fmt.Errorf("résumé: %s %d: %s", who, resp.StatusCode, strings.TrimSpace(string(b)))
}
var out summarizeResp
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
raw, err := io.ReadAll(resp.Body)
if err != nil {
return "", err
}
var out summarizeResp
if err := json.Unmarshal(raw, &out); err != nil {
return "", err
}
// Télémétrie de la compaction : réponse non streamée, timings et usage sont
// à la racine. Lue à part et sans échec possible (perfWire) : un compteur
// mal typé ne doit jamais faire échouer la compaction.
perfRecord(perfRecFromWire(perfTag{kind: perfCompact, conv: perfTagOf(ctx).conv}, decodePerfWire(raw)), nil)
if len(out.Choices) == 0 {
return "", fmt.Errorf("résumé: réponse vide")
}
+3
View File
@@ -738,6 +738,9 @@ func (c *Conversation) generate(ctx context.Context, caps Caps, temperature floa
go sendPushToAll("Loki", "Réponse prête · "+fmtDurFR(time.Since(turnStart)))
}
}()
// 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))
// Snapshot de la vue modèle.
c.mu.Lock()
+1 -1
View File
@@ -113,7 +113,7 @@ func toolSubagent(ctx context.Context, args map[string]any, parent Caps) string
// boucle (OpenFox l'a relevé). Le TEMP du preset, s'il est posé, l'emporte.
// Seul le texte écrit APRÈS le dernier outil revient au builder : la
// narration entre deux lectures n'est pas le rapport.
if _, err := runChat(ctx, msgs, subagentTemperature, caps, func(ev StreamEvent) bool {
if _, err := runChat(withPerfKind(ctx, perfSubagent), msgs, subagentTemperature, caps, func(ev StreamEvent) bool {
if ev.ToolUsed != nil {
out.Reset()
}
+1 -1
View File
@@ -115,7 +115,7 @@ func (c *Conversation) runVerifyPass(ctx context.Context, caps Caps, temperature
// meurt dans cette passe. Seuls les deltas d'affichage sortent. Son état
// dans le moteur aussi : effacé à la fin de la passe (llm_slots.go).
defer engineSideJob()()
_, _ = runChat(ctx, msgs, temperature, vcaps, func(ev StreamEvent) bool {
_, _ = runChat(withPerfKind(ctx, perfVerify), msgs, temperature, vcaps, func(ev StreamEvent) bool {
c.forwardStream(ev, epoch)
return true
})
+4
View File
@@ -143,6 +143,10 @@ func runBench(nPrompt, nPredict int) (*benchResult, error) {
PredictedN: t.PredictedN, PredictedMs: t.PredictedMs, PredictedPerSec: t.PredictedPerSec,
Elapsed: elapsed,
}
// Trace dans la télémétrie : le bench a pris le slot, la perte de cache du
// tour suivant lui revient (perf_log.go). Cache inconnu : cache_prompt:false.
perfRecord(perfRec{Kind: perfBench, Complete: true, Total: t.PromptN, New: t.PromptN,
PPms: t.PromptMs, PPtps: t.PromptPerSecond, Gen: t.PredictedN, TGtps: t.PredictedPerSec}, nil)
saveLastBench(res)
saveBenchForActivePreset(res)
return res, nil
+65 -14
View File
@@ -731,6 +731,20 @@ type StatsEvent struct {
// Conseil ponctuel quand le cache de prompts n'a pas tenu après un travail
// annexe (voir noteEnginePrompt). Vide la plupart du temps.
CacheHint string `json:"cache_hint,omitempty"`
// Télémétrie d'affichage (perf_log.go), jamais relue pour décider quoi que ce
// soit. Pointeurs : nil = le moteur ne l'a pas dit, ce qui n'est pas 0 — un
// cache_n à 0 (tout recalculé) est précisément ce qu'il faut pouvoir montrer.
CacheTokens *int `json:"cache_tokens,omitempty"`
DraftN *int `json:"draft_n,omitempty"`
DraftAccepted *int `json:"draft_accepted,omitempty"`
TTFTms *int64 `json:"ttft_ms,omitempty"`
// Posés sur le dernier événement de la complétion seulement, une fois
// celle-ci rangée : jetons du tour précédent qu'il a fallu recalculer, et
// s'il faut le signaler (seuil relevé sur un modèle hybride).
Lost *int `json:"lost,omitempty"`
LostAfter string `json:"lost_after,omitempty"`
LostAlert bool `json:"lost_alert,omitempty"`
Kind string `json:"kind,omitempty"`
}
// ChatCallback receives stream events. Return false to abort the stream.
@@ -758,13 +772,14 @@ type streamChunk struct {
PredictedPerSec float64 `json:"predicted_per_second"`
} `json:"timings"`
// Chunk final (include_usage) : taille totale du prompt, hors choices.
// ⚠️ Rien de plus ici : un champ typé de travers fait échouer le décodage
// du chunk ENTIER, qui est sauté — texte final, finish_reason et usage
// compris. Les compteurs de télémétrie (cache_n, draft_n, cached_tokens) se
// lisent à part, sans jamais échouer (perfWire, perf_log.go).
Usage *struct {
PromptTokens int `json:"prompt_tokens"`
CompletionTokens int `json:"completion_tokens"`
TotalTokens int `json:"total_tokens"`
PromptTokensDetails *struct {
CachedTokens int `json:"cached_tokens"`
} `json:"prompt_tokens_details"`
PromptTokens int `json:"prompt_tokens"`
CompletionTokens int `json:"completion_tokens"`
TotalTokens int `json:"total_tokens"`
} `json:"usage"`
}
@@ -908,6 +923,9 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
// bubble works regardless of backend. The ik_llama.cpp fork already sends
// reasoning_content, in which case we leave content untouched.
chatCfg := ReadConfig()
// Nature et conversation de ces complétions, pour la télémétrie seulement
// (perf_log.go) : rien de tout ça ne part dans la requête.
ptag, perfIter := perfTagOf(ctx), 0
reasoningOn := reasoningActive(chatCfg["REASONING"])
// Intensité du raisonnement, passée telle quelle au gabarit du modèle. Vide
// = on n'envoie rien. Réglage par preset, donc de fait par modèle.
@@ -1005,13 +1023,16 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
logBudget(toolRuns, budget, budgetNudges)
messages = append(messages, Message{Role: "user", Content: msg})
}
// Normalisé juste avant l'envoi : un seul système, en tête. Les gabarits
// stricts (Qwen3.x) refusent un système ailleurs qu'en position 0.
// Gardé à part : la télémétrie compare ces messages-là d'une requête à
// l'autre (perfPrefix).
sent := normalizeSystemMessages(messages)
payload := map[string]any{
"model": ep.Model,
// Normalisé juste avant l'envoi : un seul système, en tête. Les gabarits
// stricts (Qwen3.x) refusent un système ailleurs qu'en position 0.
// Les images de l'historique y sont rangées par référence
// (chat_images.go) : on remet leurs octets juste avant l'envoi.
"messages": expandImageRefs(normalizeSystemMessages(messages)),
"messages": expandImageRefs(sent),
"stream": true,
"temperature": temperature,
// include_usage → chunk final avec `usage.prompt_tokens` = taille TOTALE
@@ -1250,6 +1271,13 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
// build llama.cpp (MTP/spéculatif), a `choices:[]` — on les traite AVANT le
// garde de choices, sinon `gen_tokens`/`gen_per_second` (decode) sont jetés
// et l'UI retombe à « 0 tok/s » à la fin de la génération.
// Deuxième lecture, tolérante, pour la télémétrie : seulement sur les
// chunks qui portent timings ou usage, pas sur chaque jeton.
var pw perfWire
if chunk.Timings != nil || chunk.Usage != nil {
pw = decodePerfWire([]byte(data))
stats.applyPerf(pw)
}
if chunk.Timings != nil {
sawTimings = true
stats.PromptTokens = chunk.Timings.PromptN
@@ -1268,11 +1296,8 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
stats.PromptTokensTotal = chunk.Usage.PromptTokens
s := stats
if !ep.External {
cached, has := 0, chunk.Usage.PromptTokensDetails != nil
if has {
cached = chunk.Usage.PromptTokensDetails.CachedTokens
}
s.CacheHint = noteEnginePrompt(chunk.Usage.PromptTokens, cached, has)
cached := pw.cachedTokens()
s.CacheHint = noteEnginePrompt(chunk.Usage.PromptTokens, cached.int(), cached.ok)
}
scb(StreamEvent{Stats: &s})
}
@@ -1490,6 +1515,12 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
scanErr := sc.Err()
resp.Body.Close()
endReq()
// Premier jeton, de quelque nature qu'il soit (raisonnement, texte ou
// appel d'outil). tReq est repris à chaque tentative.
if !tFirst.IsZero() {
ms := tFirst.Sub(tReq).Milliseconds()
stats.TTFTms = &ms
}
if !sawTimings && !tFirst.IsZero() {
// Décodage mesuré ici ; le 1er token tombe dans la lecture du prompt.
// Pas de débit de lecture : le serveur a pu réutiliser une partie du
@@ -1507,6 +1538,26 @@ func runChat(ctx context.Context, messages []Message, temperature float64, caps
s := stats
cb(StreamEvent{Stats: &s})
}
// Télémétrie (perf_log.go). Une complétion coupée avant son chunk final
// (arrêt, appel écrit en texte, refus anticipé, flux cassé) n'a ni cache
// ni total fiables : rangée comme incomplète, sans calcul de perte.
complete := !aborted && (scanErr == nil || finishReason != "") && (sawTimings || stats.PromptTokensTotal > 0)
var prefix []uint64
if complete {
prefix = perfPrefix(sent)
}
rec := perfRecord(perfRecFromStats(ptag, perfIter, complete, stats), prefix)
perfIter++
if complete {
// Dernier événement de la complétion, copie COMPLÈTE des stats : rejoué
// depuis le journal, il ne fait que repeindre la même ligne.
s := stats
s.Kind, s.Lost, s.LostAfter = rec.Kind, rec.Lost, rec.LostAfter
if (rec.Lost != nil && *rec.Lost > 0) || (rec.Cached != nil && *rec.Cached == 0) {
s.LostAlert = perfAlert(rec, perfLostAlertAt(chatCfg))
}
cb(StreamEvent{Stats: &s})
}
if aborted {
return extra, nil
}
+5
View File
@@ -116,6 +116,11 @@ func oaiHandler() http.Handler {
// n'efface pas le slot pendant ce temps. Différé : ReverseProxy panique
// (ErrAbortHandler) quand le client coupe en plein flux.
defer engineRequestStart()()
// Complétion d'un client externe : elle prend le slot sous le nez de
// la conversation. Notée pour nommer la perte de cache qui suit.
if r.Method == http.MethodPost && strings.HasPrefix(p, "/v1/") {
perfNoteForeign()
}
lp.ServeHTTP(w, r)
return
}
+568
View File
@@ -0,0 +1,568 @@
package loki
import (
"context"
"encoding/json"
"fmt"
"hash/fnv"
"math"
"net/http"
"os"
"sort"
"strconv"
"strings"
"sync"
"time"
)
// Télémétrie par complétion : combien du prompt venait du cache, combien a été
// recalculé, à quelle vitesse, et ce que le décodage spéculatif a gardé.
//
// Rien de tout ça ne retourne au moteur : aucun champ ajouté aux requêtes, le
// comptage du contexte (usage.prompt_tokens → CtxUsed → compaction) reste
// celui d'avant. On LIT ce que llama.cpp envoie déjà sur son chunk final, on le
// garde en mémoire (anneau borné, aucun texte de message) et on l'expose en
// agrégats sur /api/perf/summary. Sans ces chiffres, impossible de savoir si une
// ligne volatile s'est glissée dans le prompt, si une mise à jour du moteur a
// cassé les points de reprise, ou si l'acceptation du MTP s'est effondrée.
// Nature d'une complétion. Elle voyage dans le contexte (withPerf) et jamais
// dans Caps : Caps est une liste de capacités, recopiée par valeur partout.
const (
perfMain = "main" // tour de la conversation (et ses étapes d'outils)
perfSubagent = "subagent" // sous-agent du mode Code
perfVerify = "verify" // passe de vérification du mode Code
perfTask = "task" // tâche planifiée
perfCompact = "compact" // résumé de compaction
perfBench = "bench" // mesure du moteur (cache_prompt:false)
perfForeign = "foreign" // client externe passé par /v1 : vu, pas mesuré
)
type perfCtxKey struct{}
type perfTag struct{ kind, conv string }
// withPerf étiquette les complétions lancées sous ctx.
func withPerf(ctx context.Context, kind, conv string) context.Context {
return context.WithValue(ctx, perfCtxKey{}, perfTag{kind: kind, conv: conv})
}
// withPerfKind change la nature en gardant la conversation : un sous-agent ou
// une vérification reste rattaché à la discussion qui l'a lancé.
func withPerfKind(ctx context.Context, kind string) context.Context {
t := perfTagOf(ctx)
t.kind = kind
return context.WithValue(ctx, perfCtxKey{}, t)
}
func perfTagOf(ctx context.Context) perfTag {
t, _ := ctx.Value(perfCtxKey{}).(perfTag)
if t.kind == "" {
t.kind = perfMain
}
return t
}
// perfNum : un nombre de télémétrie, décodé sans JAMAIS échouer. Un fork ou une
// API tierce qui l'envoie en flottant, en chaîne ou pas du tout le laisse
// simplement « inconnu » : une valeur d'affichage ne doit pas pouvoir jeter le
// chunk final (texte, finish_reason, usage) ni faire échouer une compaction.
type perfNum struct {
f float64
ok bool
}
func (n *perfNum) UnmarshalJSON(b []byte) error {
s := strings.Trim(strings.TrimSpace(string(b)), `"`)
if f, err := strconv.ParseFloat(s, 64); err == nil && !math.IsNaN(f) && !math.IsInf(f, 0) && f >= 0 && f < 1e12 {
*n = perfNum{f: f, ok: true}
}
return nil
}
func (n perfNum) int() int { return int(n.f) }
func (n perfNum) ptr() *int {
if !n.ok {
return nil
}
v := int(n.f)
return &v
}
// perfWire : les champs de télémétrie du chunk final, lus dans un DEUXIÈME
// décodage, séparé de streamChunk et dont l'erreur est ignorée. Les compteurs
// cache_n/draft_n n'existent pas sur tous les moteurs (draft_n n'apparaît
// qu'avec un brouillon actif) : absents, ils restent inconnus, jamais 0.
type perfWire struct {
Timings *struct {
PromptN perfNum `json:"prompt_n"`
PromptMs perfNum `json:"prompt_ms"`
PromptPerSecond perfNum `json:"prompt_per_second"`
PredictedN perfNum `json:"predicted_n"`
PredictedPerSec perfNum `json:"predicted_per_second"`
CacheN perfNum `json:"cache_n"`
DraftN perfNum `json:"draft_n"`
DraftNAccepted perfNum `json:"draft_n_accepted"`
} `json:"timings"`
Usage *struct {
PromptTokens perfNum `json:"prompt_tokens"`
CompletionTokens perfNum `json:"completion_tokens"`
PromptTokensDetails *struct {
CachedTokens perfNum `json:"cached_tokens"`
} `json:"prompt_tokens_details"`
} `json:"usage"`
}
func decodePerfWire(data []byte) perfWire {
var w perfWire
_ = json.Unmarshal(data, &w) // ce qui a pu être lu l'est ; le reste reste inconnu
return w
}
// cachedTokens : usage.prompt_tokens_details.cached_tokens (format OpenAI).
func (w perfWire) cachedTokens() perfNum {
if w.Usage == nil || w.Usage.PromptTokensDetails == nil {
return perfNum{}
}
return w.Usage.PromptTokensDetails.CachedTokens
}
// applyPerf reporte la télémétrie d'un chunk sur les stats. timings.cache_n a
// la priorité ; cached_tokens ne fait que combler un cache encore inconnu (les
// deux viennent du même compteur côté llama.cpp, mais seul le second existe
// sur une API tierce).
func (s *StatsEvent) applyPerf(w perfWire) {
if t := w.Timings; t != nil {
if p := t.CacheN.ptr(); p != nil {
s.CacheTokens = p
}
if p := t.DraftN.ptr(); p != nil {
s.DraftN = p
}
if p := t.DraftNAccepted.ptr(); p != nil {
s.DraftAccepted = p
}
}
if s.CacheTokens == nil {
s.CacheTokens = w.cachedTokens().ptr()
}
}
// perfRec : une complétion. Des identifiants et des compteurs, jamais de texte.
type perfRec struct {
At time.Time `json:"at"`
Seq uint64 `json:"seq"`
Kind string `json:"kind"`
Conv string `json:"conv,omitempty"`
Iter int `json:"iter"`
Complete bool `json:"complete"`
// Total : usage.prompt_tokens, la taille du prompt entier. New : prompt_n,
// les jetons réellement recalculés. Cached : cache_n, absent si inconnu.
Total int `json:"total,omitempty"`
Cached *int `json:"cached,omitempty"`
New int `json:"new,omitempty"`
PPms float64 `json:"pp_ms,omitempty"`
PPtps float64 `json:"pp_tps,omitempty"`
Gen int `json:"gen,omitempty"`
TGtps float64 `json:"tg_tps,omitempty"`
DraftN *int `json:"draft_n,omitempty"`
DraftAcc *int `json:"draft_accepted,omitempty"`
TTFTms *int64 `json:"ttft_ms,omitempty"`
Lost *int `json:"lost,omitempty"`
LostAfter string `json:"lost_after,omitempty"`
}
// perfRecFromStats : la complétion telle que runChat l'a vue.
func perfRecFromStats(t perfTag, iter int, complete bool, s StatsEvent) perfRec {
return perfRec{
Kind: t.kind, Conv: t.conv, Iter: iter, Complete: complete,
Total: s.PromptTokensTotal, Cached: s.CacheTokens, New: s.PromptTokens,
PPms: s.PromptMs, PPtps: s.PromptPerSecond, Gen: s.GenTokens, TGtps: s.GenPerSecond,
DraftN: s.DraftN, DraftAcc: s.DraftAccepted, TTFTms: s.TTFTms,
}
}
// perfRecFromWire : une réponse NON streamée (compaction), lue d'un bloc.
func perfRecFromWire(t perfTag, w perfWire) perfRec {
r := perfRec{Kind: t.kind, Conv: t.conv}
if tm := w.Timings; tm != nil {
r.New, r.PPms, r.PPtps = tm.PromptN.int(), tm.PromptMs.f, tm.PromptPerSecond.f
r.Gen, r.TGtps = tm.PredictedN.int(), tm.PredictedPerSec.f
r.Cached, r.DraftN, r.DraftAcc = tm.CacheN.ptr(), tm.DraftN.ptr(), tm.DraftNAccepted.ptr()
}
if u := w.Usage; u != nil {
r.Total = u.PromptTokens.int()
if r.Gen == 0 {
r.Gen = u.CompletionTokens.int()
}
}
if r.Cached == nil {
r.Cached = w.cachedTokens().ptr()
}
r.Complete = r.Total > 0 || w.Timings != nil
return r
}
// perfPrefix : empreintes cumulées des messages envoyés (l'empreinte i couvre
// les messages 0..i). C'est ce qui dit qu'une requête PROLONGE la précédente
// de la même conversation — l'identifiant seul ne suffit pas : une compaction,
// une édition ou une autre tâche dans la même discussion repartent d'ailleurs.
func perfPrefix(msgs []Message) []uint64 {
h := fnv.New64a()
out := make([]uint64, len(msgs))
for i, m := range msgs {
b, _ := json.Marshal(m)
h.Write(b)
h.Write([]byte{0})
out[i] = h.Sum64()
}
return out
}
const (
perfRingMax = 5000 // ~1 Mio en mémoire au pire
perfConvMax = 128 // conversations suivies pour le calcul de la perte
perfBetween = 500 // entrées remontées au plus pour nommer ce qui s'est intercalé
)
// perfConvState : la dernière complétion terminée d'une conversation (et d'une
// nature : le fil d'un sous-agent n'est pas celui de son builder).
type perfConvState struct {
seq uint64
n int
hash uint64
total int
}
type perfStore struct {
mu sync.Mutex
seq uint64
ring []perfRec
head int // prochaine case à écraser quand l'anneau est plein
convs map[string]perfConvState
}
var perfLog = &perfStore{}
// record range la complétion et calcule sa perte de cache :
//
// lost = max(0, total précédent − cache_n)
//
// seulement quand le cache est connu, que la complétion est allée au bout, et
// que ses messages PROLONGENT strictement ceux de la précédente complétion de
// la même conversation et de la même nature. Ailleurs (nouvelle discussion,
// compaction, flux coupé) la perte reste inconnue plutôt que fausse.
func (p *perfStore) record(r perfRec, prefix []uint64) perfRec {
p.mu.Lock()
defer p.mu.Unlock()
p.seq++
r.Seq = p.seq
if r.At.IsZero() {
r.At = time.Now()
}
key := r.Kind + "\x00" + r.Conv
if p.convs == nil {
p.convs = map[string]perfConvState{}
}
if r.Complete && len(prefix) > 0 {
prev, had := p.convs[key]
if had && r.Cached != nil && prev.total > 0 && prev.n > 0 && len(prefix) > prev.n && prefix[prev.n-1] == prev.hash {
lost := max(0, prev.total-*r.Cached)
r.Lost = &lost
if lost > 0 {
r.LostAfter = p.between(prev.seq, key)
}
}
if r.Total > 0 {
p.convs[key] = perfConvState{seq: r.Seq, n: len(prefix), hash: prefix[len(prefix)-1], total: r.Total}
p.evictConvs()
}
}
if len(p.ring) < perfRingMax {
p.ring = append(p.ring, r)
} else {
p.ring[p.head] = r
p.head = (p.head + 1) % perfRingMax
}
return r
}
// between nomme les natures des complétions passées par le moteur depuis seq
// (un sous-agent, un client externe…) : la cause probable d'une perte. Vide =
// rien vu d'autre, perte non attribuée.
func (p *perfStore) between(seq uint64, key string) string {
seen := map[string]bool{}
n := len(p.ring)
// De la plus récente à la plus ancienne. Tant que l'anneau n'est pas plein,
// head vaut 0 et la plus récente est la dernière case : même formule.
for i := 0; i < n && i < perfBetween; i++ {
e := p.ring[((p.head-1-i)%n+n)%n]
if e.Seq <= seq {
break
}
if e.Kind+"\x00"+e.Conv != key {
seen[e.Kind] = true
}
}
kinds := make([]string, 0, len(seen))
for k := range seen {
kinds = append(kinds, k)
}
sort.Strings(kinds)
return strings.Join(kinds, ",")
}
// evictConvs borne l'état par conversation : on oublie la plus ancienne.
func (p *perfStore) evictConvs() {
for len(p.convs) > perfConvMax {
oldK, oldS := "", uint64(math.MaxUint64)
for k, s := range p.convs {
if s.seq < oldS {
oldK, oldS = k, s.seq
}
}
delete(p.convs, oldK)
}
}
// snapshot : copie de l'anneau, de la plus ancienne à la plus récente.
func (p *perfStore) snapshot() []perfRec {
p.mu.Lock()
defer p.mu.Unlock()
out := make([]perfRec, 0, len(p.ring))
out = append(out, p.ring[p.head:]...)
return append(out, p.ring[:p.head]...)
}
// perfRecord range une complétion et, si LOKI_PERF_LOG est posée, l'écrit sur
// stderr. Une ligne par complétion, c'est des dizaines par tour d'agent : pas
// par défaut, les journaux Docker ont déjà celles de llama-server.
func perfRecord(r perfRec, prefix []uint64) perfRec {
r = perfLog.record(r, prefix)
if perfLogOn() {
fmt.Fprintln(os.Stderr, perfLine(r))
}
return r
}
// perfNoteForeign : une requête d'un client externe vient de passer par /v1.
// Ses chiffres ne sont pas lus (le flux est relayé tel quel), mais sa trace
// suffit à expliquer la perte de cache du tour suivant.
func perfNoteForeign() {
perfRecord(perfRec{Kind: perfForeign}, nil)
}
func perfLogOn() bool {
v := strings.ToLower(strings.TrimSpace(os.Getenv("LOKI_PERF_LOG")))
return v != "" && v != "0" && v != "false" && v != "off" && v != "no"
}
func perfLine(r perfRec) string {
opt := func(p *int) string {
if p == nil {
return "?"
}
return strconv.Itoa(*p)
}
draft := "?"
if r.DraftN != nil {
draft = opt(r.DraftAcc) + "/" + strconv.Itoa(*r.DraftN)
}
ttft := "?"
if r.TTFTms != nil {
ttft = strconv.FormatInt(*r.TTFTms, 10)
}
lost := opt(r.Lost)
if r.LostAfter != "" {
lost += "(" + r.LostAfter + ")"
}
conv := r.Conv
if conv == "" {
conv = "-"
}
return fmt.Sprintf("[perf] kind=%s conv=%s iter=%d total=%d cached=%s new=%d pp_ms=%.0f pp_tps=%.1f gen=%d tg_tps=%.1f draft=%s ttft_ms=%s lost=%s complete=%t",
r.Kind, conv, r.Iter, r.Total, opt(r.Cached), r.New, r.PPms, r.PPtps, r.Gen, r.TGtps, draft, ttft, lost, r.Complete)
}
// perfLostAlertAt : la perte au-delà de laquelle l'interface la signale. Sur un
// modèle hybride (couches récurrentes), le cache ne reprend qu'à un point de
// reprise : perdre jusqu'à un espacement de points (--checkpoint-min-step,
// 2048 posé par Loki) plus un micro-batch est l'état NORMAL, pas un incident.
// Le signaler à chaque étape noierait les vraies pertes.
func perfLostAlertAt(cfg map[string]string) int {
const base = 1000
step := 0
if n, err := strconv.Atoi(strings.TrimSpace(cfg["CKPT_MIN_STEP"])); err == nil && n > 0 {
step = n
} else if m := strings.TrimSpace(cfg["MODEL"]); m != "" {
if p, err := resolveServeModelPath(m); err == nil {
if g, err := ggufMeta(p); err == nil && ggufHybrid(&g) {
step = 2048
}
}
}
if step == 0 {
return base
}
ub := 512
if n, err := strconv.Atoi(strings.TrimSpace(cfg["UBATCH"])); err == nil && n > 0 {
ub = n
}
return max(base, step+ub)
}
// perfAlert : la complétion mérite d'être signalée — une perte au-delà du seuil,
// ou un gros prompt repris de zéro alors que le cache est connu.
func perfAlert(r perfRec, threshold int) bool {
if r.Lost != nil && *r.Lost > threshold {
return true
}
return r.Cached != nil && *r.Cached == 0 && r.Total >= 8*threshold
}
// --- /api/perf/summary -------------------------------------------------------
type perfKindSum struct {
Completions int `json:"completions"`
Incomplete int `json:"incomplete"`
// Turns : premières complétions (iter 0) — un tour utilisateur, un
// sous-agent, une passe de vérification…
Turns int `json:"turns"`
NewTokens int `json:"new_tokens"`
PrefillSec float64 `json:"prefill_sec"`
NewTokensPerTurn float64 `json:"new_tokens_per_turn,omitempty"`
PrefillSecPerTurn float64 `json:"prefill_sec_per_turn,omitempty"`
LostEvents int `json:"lost_events"`
LostTokens int `json:"lost_tokens"`
TTFTMedianMs *float64 `json:"ttft_median_ms,omitempty"`
}
type perfDepthSum struct {
Bucket string `json:"bucket"`
N int `json:"n"`
PPMedianTPS *float64 `json:"pp_median_tps,omitempty"`
TGMedianTPS *float64 `json:"tg_median_tps,omitempty"`
}
type perfSummary struct {
Entries int `json:"entries"`
Since time.Time `json:"since,omitzero"`
// CacheHitRatio : Σ cached / Σ total sur les complétions où les deux sont
// connus. Absent = moteur qui ne dit rien de son cache.
CacheHitRatio *float64 `json:"cache_hit_ratio,omitempty"`
Kinds map[string]*perfKindSum `json:"kinds"`
Depth []perfDepthSum `json:"depth"`
DraftN int `json:"draft_n"`
DraftAccepted int `json:"draft_accepted"`
DraftRate *float64 `json:"draft_rate,omitempty"`
}
var perfDepthBuckets = []struct {
name string
max int
}{{"0-8k", 8 << 10}, {"8-16k", 16 << 10}, {"16-32k", 32 << 10}, {"32k+", math.MaxInt}}
func perfMedian(v []float64) *float64 {
if len(v) == 0 {
return nil
}
sort.Float64s(v)
m := v[len(v)/2]
if len(v)%2 == 0 {
m = (v[len(v)/2-1] + v[len(v)/2]) / 2
}
return &m
}
func perfSummarize(recs []perfRec) perfSummary {
s := perfSummary{Entries: len(recs), Kinds: map[string]*perfKindSum{}}
if len(recs) > 0 {
s.Since = recs[0].At
}
var sumCached, sumTotal int
pp := make([][]float64, len(perfDepthBuckets))
tg := make([][]float64, len(perfDepthBuckets))
ttft := map[string][]float64{}
for _, r := range recs {
k := s.Kinds[r.Kind]
if k == nil {
k = &perfKindSum{}
s.Kinds[r.Kind] = k
}
if !r.Complete {
k.Incomplete++
continue
}
k.Completions++
if r.Iter == 0 {
k.Turns++
}
k.NewTokens += r.New
k.PrefillSec += r.PPms / 1000
if r.Lost != nil && *r.Lost > 0 {
k.LostEvents++
k.LostTokens += *r.Lost
}
if r.TTFTms != nil {
ttft[r.Kind] = append(ttft[r.Kind], float64(*r.TTFTms))
}
if r.Cached != nil && r.Total > 0 {
sumCached += min(*r.Cached, r.Total)
sumTotal += r.Total
}
depth := r.Total
if depth == 0 {
depth = r.New
}
for i, b := range perfDepthBuckets {
if depth < b.max {
// Débit de lecture : seulement sur un vrai prefill (quelques jetons
// nouveaux donnent des débits sans signification).
if r.PPtps > 0 && r.New >= 64 {
pp[i] = append(pp[i], r.PPtps)
}
if r.TGtps > 0 && r.Gen > 1 {
tg[i] = append(tg[i], r.TGtps)
}
break
}
}
if r.DraftN != nil {
s.DraftN += *r.DraftN
if r.DraftAcc != nil {
s.DraftAccepted += *r.DraftAcc
}
}
}
for kind, k := range s.Kinds {
if k.Turns > 0 {
k.NewTokensPerTurn = float64(k.NewTokens) / float64(k.Turns)
k.PrefillSecPerTurn = k.PrefillSec / float64(k.Turns)
}
k.TTFTMedianMs = perfMedian(ttft[kind])
}
if sumTotal > 0 {
r := float64(sumCached) / float64(sumTotal)
s.CacheHitRatio = &r
}
if s.DraftN > 0 {
r := float64(s.DraftAccepted) / float64(s.DraftN)
s.DraftRate = &r
}
for i, b := range perfDepthBuckets {
s.Depth = append(s.Depth, perfDepthSum{Bucket: b.name, N: max(len(pp[i]), len(tg[i])),
PPMedianTPS: perfMedian(pp[i]), TGMedianTPS: perfMedian(tg[i])})
}
return s
}
// handlePerfSummary : GET /api/perf/summary (derrière la clé de pilotage).
// Des agrégats seulement : aucun texte de message n'est jamais retenu.
func handlePerfSummary(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
sendJSON(w, http.StatusMethodNotAllowed, map[string]any{"error": "GET seulement"})
return
}
sendJSON(w, http.StatusOK, perfSummarize(perfLog.snapshot()))
}
+329
View File
@@ -0,0 +1,329 @@
package loki
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"net/url"
"strings"
"testing"
)
// freshPerf isole l'anneau de télémétrie d'un test.
func freshPerf(t *testing.T) {
t.Helper()
old := perfLog
perfLog = &perfStore{}
t.Cleanup(func() { perfLog = old })
}
func iptr(v int) *int { return &v }
// Chunk final d'un llama-server récent, en flux, avec MTP actif : choices
// vide, timings et usage à la racine.
const llamaFinalChunk = `{"choices":[],"created":1,"id":"x","model":"m","object":"chat.completion.chunk",` +
`"usage":{"completion_tokens":120,"prompt_tokens":41484,"total_tokens":41604,"prompt_tokens_details":{"cached_tokens":41230}},` +
`"timings":{"cache_n":41230,"prompt_n":254,"prompt_ms":310.2,"prompt_per_token":1.2,"prompt_per_second":818.8,` +
`"predicted_n":120,"predicted_ms":2400,"predicted_per_token":20,"predicted_per_second":50,"draft_n":90,"draft_n_accepted":71}}`
func TestPerfWireDecode(t *testing.T) {
cases := []struct {
name string
data string
cache, draft, accept *int
}{
{"chunk final llama.cpp", llamaFinalChunk, iptr(41230), iptr(90), iptr(71)},
// Sans brouillon, llama.cpp n'émet pas draft_n : inconnu, pas 0.
{"sans brouillon", `{"timings":{"cache_n":12,"prompt_n":3}}`, iptr(12), nil, nil},
// Tout recalculé : un 0 CONNU, qu'il faut pouvoir montrer.
{"cache vide", `{"timings":{"cache_n":0,"prompt_n":3000}}`, iptr(0), nil, nil},
// API tierce : seulement le format OpenAI.
{"usage seul", `{"usage":{"prompt_tokens":900,"prompt_tokens_details":{"cached_tokens":512}}}`, iptr(512), nil, nil},
{"moteur muet", `{"usage":{"prompt_tokens":900}}`, nil, nil, nil},
// Types de travers (fork, passerelle) : lu si c'est un nombre, sinon inconnu.
{"types de travers", `{"timings":{"cache_n":"abc","draft_n":"12","draft_n_accepted":7.0}}`, nil, iptr(12), iptr(7)},
{"structure de travers", `{"timings":"nope","usage":{"prompt_tokens_details":{"cached_tokens":5.9}}}`, iptr(5), nil, nil},
}
eq := func(a, b *int) bool { return (a == nil) == (b == nil) && (a == nil || *a == *b) }
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
var s StatsEvent
s.applyPerf(decodePerfWire([]byte(c.data)))
if !eq(s.CacheTokens, c.cache) || !eq(s.DraftN, c.draft) || !eq(s.DraftAccepted, c.accept) {
t.Fatalf("cache=%v draft=%v accepted=%v", s.CacheTokens, s.DraftN, s.DraftAccepted)
}
})
}
}
// Le chunk final peut aussi porter le dernier texte et le finish_reason (forks,
// API compatibles). Un compteur de télémétrie mal typé ne doit pas le faire
// sauter : avant, cached_tokens strictement entier jetait tout le chunk.
func TestStreamChunkIgnoresOddTelemetry(t *testing.T) {
data := `{"choices":[{"delta":{"content":"fin"},"finish_reason":"stop"}],` +
`"usage":{"prompt_tokens":500,"completion_tokens":3,"prompt_tokens_details":{"cached_tokens":480.5}},` +
`"timings":{"prompt_n":20,"predicted_n":3,"draft_n":"x","cache_n":480.0}}`
var c streamChunk
if err := json.Unmarshal([]byte(data), &c); err != nil {
t.Fatalf("chunk jeté : %v", err)
}
if len(c.Choices) != 1 || c.Choices[0].Delta.Content != "fin" || c.Choices[0].FinishReason != "stop" || c.Usage.PromptTokens != 500 {
t.Fatalf("chunk mal lu : %+v", c)
}
var s StatsEvent
s.applyPerf(decodePerfWire([]byte(data)))
if s.CacheTokens == nil || *s.CacheTokens != 480 || s.DraftN != nil {
t.Fatalf("télémétrie : cache=%v draft=%v", s.CacheTokens, s.DraftN)
}
}
// De bout en bout : le dernier événement Stats porte la télémétrie, le texte et
// le comptage du contexte sont intacts, et rien n'est ajouté à la requête.
func TestRunChatPerfFinalStats(t *testing.T) {
testHome(t)
freshPerf(t)
var payload map[string]any
body := sseChunk("bonjour") +
`data: {"choices":[{"delta":{"content":" toi"},"finish_reason":"stop"}],"usage":{"prompt_tokens":500,"completion_tokens":2,"prompt_tokens_details":{"cached_tokens":480.5}},"timings":{"prompt_n":20,"prompt_ms":10,"prompt_per_second":2000,"predicted_n":2,"predicted_ms":40,"predicted_per_second":50,"draft_n":"x"}}` + "\n\n" +
"data: [DONE]\n\n"
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
_ = json.NewDecoder(r.Body).Decode(&payload)
w.Header().Set("Content-Type", "text/event-stream")
_, _ = w.Write([]byte(body))
}))
t.Cleanup(srv.Close)
u, _ := url.Parse(srv.URL)
if err := SetConfigKey("PORT", u.Port()); err != nil {
t.Fatal(err)
}
var content strings.Builder
var last *StatsEvent
ctx := withPerf(context.Background(), perfMain, "c1")
if _, err := runChat(ctx, []Message{{Role: "user", Content: "salut"}}, 0.7, Caps{}, func(ev StreamEvent) bool {
if ev.Content != "" {
content.WriteString(ev.Content)
}
if ev.Stats != nil {
last = ev.Stats
}
return true
}); err != nil {
t.Fatal(err)
}
if content.String() != "bonjour toi" {
t.Fatalf("texte : %q", content.String())
}
if last == nil || last.PromptTokensTotal != 500 || last.Kind != perfMain || last.CacheTokens == nil || *last.CacheTokens != 480 || last.TTFTms == nil {
t.Fatalf("dernier Stats : %+v", last)
}
for _, k := range []string{"timings_per_token", "return_progress", "n_probs", "id_slot", "cache_prompt"} {
if _, ok := payload[k]; ok {
t.Fatalf("la télémétrie a ajouté %q à la requête", k)
}
}
recs := perfLog.snapshot()
if len(recs) != 1 || !recs[0].Complete || recs[0].Conv != "c1" || recs[0].Total != 500 || recs[0].New != 20 {
t.Fatalf("enregistrement : %+v", recs)
}
}
// Un flux coupé n'a ni cache ni total fiables : rangé incomplet, sans perte.
func TestRunChatPerfCutStreamIncomplete(t *testing.T) {
testHome(t)
freshPerf(t)
port := sseCuttingServer(t, sseChunk("début"), true)
if err := SetConfigKey("PORT", port); err != nil {
t.Fatal(err)
}
_, _ = runChat(context.Background(), []Message{{Role: "user", Content: "x"}}, 0.7, Caps{}, func(StreamEvent) bool { return true })
recs := perfLog.snapshot()
if len(recs) != 1 || recs[0].Complete || recs[0].Lost != nil {
t.Fatalf("flux coupé : %+v", recs)
}
}
func TestPerfRecordLost(t *testing.T) {
a := Message{Role: "user", Content: "a"}
b := Message{Role: "assistant", Content: "b"}
c := Message{Role: "tool", Content: "c"}
z := Message{Role: "user", Content: "z"}
type step struct {
rec perfRec
msgs []Message
}
main := func(total int, cached *int) perfRec {
return perfRec{Kind: perfMain, Conv: "c1", Complete: true, Total: total, Cached: cached}
}
cases := []struct {
name string
steps []step
lost *int
lostAfter string
}{
{"prolonge, cache intact", []step{{main(1000, iptr(0)), []Message{a}}, {main(1200, iptr(1000)), []Message{a, b, c}}}, iptr(0), ""},
{"prolonge, cache perdu", []step{{main(1000, iptr(0)), []Message{a}}, {main(1200, iptr(300)), []Message{a, b, c}}}, iptr(700), ""},
{"diverge", []step{{main(1000, iptr(0)), []Message{a, b}}, {main(1200, iptr(10)), []Message{z, b, c}}}, nil, ""},
{"même liste, pas une suite", []step{{main(1000, iptr(0)), []Message{a, b}}, {main(1000, iptr(10)), []Message{a, b}}}, nil, ""},
{"cache inconnu", []step{{main(1000, iptr(0)), []Message{a}}, {main(1200, nil), []Message{a, b}}}, nil, ""},
{"complétion coupée", []step{{main(1000, iptr(0)), []Message{a}}, {perfRec{Kind: perfMain, Conv: "c1", Cached: iptr(0)}, []Message{a, b}}}, nil, ""},
{"autre discussion", []step{{main(1000, iptr(0)), []Message{a}}, {perfRec{Kind: perfMain, Conv: "c2", Complete: true, Total: 1200, Cached: iptr(0)}, []Message{a, b}}}, nil, ""},
{"un sous-agent ne se compare pas au builder", []step{
{main(1000, iptr(0)), []Message{a}},
{perfRec{Kind: perfSubagent, Conv: "c1", Complete: true, Total: 1200, Cached: iptr(0)}, []Message{a, b}},
}, nil, ""},
{"perte après sous-agent et client externe", []step{
{main(1000, iptr(0)), []Message{a}},
{perfRec{Kind: perfSubagent, Conv: "c1", Complete: true, Total: 300, Cached: iptr(0)}, []Message{z}},
{perfRec{Kind: perfForeign}, nil},
{main(1300, iptr(200)), []Message{a, b, c}},
}, iptr(800), "foreign,subagent"},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
freshPerf(t)
var last perfRec
for _, s := range tc.steps {
var prefix []uint64
if s.msgs != nil {
prefix = perfPrefix(s.msgs)
}
last = perfLog.record(s.rec, prefix)
}
if (last.Lost == nil) != (tc.lost == nil) || (last.Lost != nil && *last.Lost != *tc.lost) {
t.Fatalf("lost = %v, attendu %v", last.Lost, tc.lost)
}
if last.LostAfter != tc.lostAfter {
t.Fatalf("lost_after = %q, attendu %q", last.LostAfter, tc.lostAfter)
}
})
}
}
func TestPerfRingBounded(t *testing.T) {
freshPerf(t)
for i := 0; i < perfRingMax+10; i++ {
perfLog.record(perfRec{Kind: perfMain, Conv: string(rune('a' + i%200)), Complete: true, Total: 10}, []uint64{uint64(i)})
}
recs := perfLog.snapshot()
if len(recs) != perfRingMax {
t.Fatalf("anneau : %d entrées", len(recs))
}
for i := 1; i < len(recs); i++ {
if recs[i].Seq != recs[i-1].Seq+1 {
t.Fatalf("ordre rompu en %d : %d puis %d", i, recs[i-1].Seq, recs[i].Seq)
}
}
if recs[len(recs)-1].Seq != perfRingMax+10 {
t.Fatalf("dernière entrée : %d", recs[len(recs)-1].Seq)
}
if len(perfLog.convs) > perfConvMax {
t.Fatalf("état par conversation non borné : %d", len(perfLog.convs))
}
}
func TestPerfSummarize(t *testing.T) {
recs := []perfRec{
{Kind: perfMain, Iter: 0, Complete: true, Total: 4000, Cached: iptr(3000), New: 1000, PPms: 2000, PPtps: 500, Gen: 50, TGtps: 40, DraftN: iptr(40), DraftAcc: iptr(30), TTFTms: ptr64(2100)},
{Kind: perfMain, Iter: 1, Complete: true, Total: 20000, Cached: iptr(19000), New: 1000, PPms: 1000, PPtps: 1000, Gen: 10, TGtps: 30, Lost: iptr(500), TTFTms: ptr64(1100)},
{Kind: perfSubagent, Iter: 0, Complete: true, Total: 2000, New: 2000, PPms: 4000, PPtps: 500},
{Kind: perfMain, Iter: 0},
{Kind: perfForeign},
}
s := perfSummarize(recs)
if s.Entries != 5 || s.CacheHitRatio == nil || *s.CacheHitRatio != 22000.0/24000.0 {
t.Fatalf("ratio : %+v", s.CacheHitRatio)
}
m := s.Kinds[perfMain]
if m.Completions != 2 || m.Incomplete != 1 || m.Turns != 1 || m.NewTokensPerTurn != 2000 || m.PrefillSecPerTurn != 3 || m.LostEvents != 1 || m.LostTokens != 500 {
t.Fatalf("main : %+v", m)
}
if m.TTFTMedianMs == nil || *m.TTFTMedianMs != 1600 {
t.Fatalf("ttft médian : %v", m.TTFTMedianMs)
}
if s.Kinds[perfForeign].Incomplete != 1 {
t.Fatalf("foreign : %+v", s.Kinds[perfForeign])
}
if s.DraftRate == nil || *s.DraftRate != 0.75 {
t.Fatalf("acceptation : %v", s.DraftRate)
}
if len(s.Depth) != 4 || s.Depth[0].PPMedianTPS == nil || *s.Depth[0].PPMedianTPS != 500 || s.Depth[2].TGMedianTPS == nil || *s.Depth[2].TGMedianTPS != 30 || s.Depth[3].N != 0 {
t.Fatalf("profondeurs : %+v", s.Depth)
}
// Aucune donnée : pas de ratio inventé.
if e := perfSummarize(nil); e.CacheHitRatio != nil || e.DraftRate != nil {
t.Fatalf("vide : %+v", e)
}
}
func ptr64(v int64) *int64 { return &v }
func TestPerfLostAlert(t *testing.T) {
cases := []struct {
name string
cfg map[string]string
want int
}{
{"modèle classique", map[string]string{}, 1000},
{"espacement posé", map[string]string{"CKPT_MIN_STEP": "2048"}, 2560},
{"espacement et micro-batch", map[string]string{"CKPT_MIN_STEP": "4096", "UBATCH": "1024"}, 5120},
{"espacement illisible", map[string]string{"CKPT_MIN_STEP": "beaucoup"}, 1000},
}
for _, c := range cases {
if got := perfLostAlertAt(c.cfg); got != c.want {
t.Errorf("%s : seuil %d, attendu %d", c.name, got, c.want)
}
}
if perfAlert(perfRec{Lost: iptr(900)}, 1000) || !perfAlert(perfRec{Lost: iptr(1500)}, 1000) {
t.Fatal("seuil de perte")
}
if perfAlert(perfRec{Lost: iptr(2000)}, 2560) {
t.Fatal("hybride : une perte sous l'espacement des points de reprise est normale")
}
if !perfAlert(perfRec{Cached: iptr(0), Total: 9000}, 1000) || perfAlert(perfRec{Cached: iptr(0), Total: 500}, 1000) || perfAlert(perfRec{Total: 9000}, 1000) {
t.Fatal("gros prompt sans cache")
}
}
// La compaction n'est pas streamée : sa télémétrie vient de la racine de la
// réponse, et un compteur mal typé ne la fait pas échouer.
func TestSummarizePerfLenient(t *testing.T) {
testHome(t)
freshPerf(t)
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
_, _ = w.Write([]byte(`{"choices":[{"message":{"content":"- résumé"}}],` +
`"usage":{"prompt_tokens":3000,"completion_tokens":40,"prompt_tokens_details":{"cached_tokens":"beaucoup"}},` +
`"timings":{"cache_n":2900.0,"prompt_n":100,"prompt_ms":50,"prompt_per_second":2000,"predicted_n":40,"predicted_per_second":30}}`))
}))
t.Cleanup(srv.Close)
u, _ := url.Parse(srv.URL)
if err := SetConfigKey("PORT", u.Port()); err != nil {
t.Fatal(err)
}
got, err := summarizeTranscriptFor(withPerf(context.Background(), perfMain, "c9"), "transcript", false)
if err != nil || got != "- résumé" {
t.Fatalf("résumé : %q, %v", got, err)
}
recs := perfLog.snapshot()
if len(recs) != 1 || recs[0].Kind != perfCompact || recs[0].Conv != "c9" || !recs[0].Complete ||
recs[0].Cached == nil || *recs[0].Cached != 2900 || recs[0].Total != 3000 || recs[0].New != 100 {
t.Fatalf("enregistrement : %+v", recs)
}
}
func TestPerfLine(t *testing.T) {
l := perfLine(perfRec{Kind: perfMain, Conv: "c1", Total: 41484, Cached: iptr(41230), New: 254, DraftN: iptr(90), DraftAcc: iptr(71), Lost: iptr(0), Complete: true})
for _, want := range []string{"kind=main", "conv=c1", "cached=41230", "new=254", "draft=71/90", "ttft_ms=?", "lost=0"} {
if !strings.Contains(l, want) {
t.Fatalf("%q absent de %q", want, l)
}
}
t.Setenv("LOKI_PERF_LOG", "")
if perfLogOn() {
t.Fatal("ligne [perf] active par défaut")
}
t.Setenv("LOKI_PERF_LOG", "1")
if !perfLogOn() {
t.Fatal("LOKI_PERF_LOG=1 ignorée")
}
}
+1 -1
View File
@@ -93,7 +93,7 @@ func (c *Conversation) RunAutonomous(ctx context.Context, taskID, taskName, prom
// La tâche a pris le slot de la conversation : effacé une fois la tâche
// finie, verrou de génération encore tenu (llm_slots.go).
defer engineSideJob()()
_, err := runChat(ctx, InjectSkills(final, caps), temperature, caps, func(ev StreamEvent) bool {
_, err := runChat(withPerf(ctx, perfTask, "task:"+taskID), InjectSkills(final, caps), temperature, caps, func(ev StreamEvent) bool {
if ev.Content != "" {
content.WriteString(ev.Content)
}
+22 -3
View File
@@ -1151,6 +1151,9 @@ h1{color:var(--text);font-weight:600;letter-spacing:-.01em}
#chat .msg .statline{margin-top:10px;font-family:var(--mono);font-size:10.5px;letter-spacing:.02em;
color:var(--dim);opacity:.75}
#chat .msg .statline:empty{display:none}
/* Cache de prompts perdu (gros recalcul) : la vitesse passe en ambre, la
ligne reste discrète mais ne se lit plus comme un tour ordinaire. */
#chat .msg .statline.cache-miss>.spd{color:var(--warn)}
/* « Masquer la vitesse de génération » ne masque QUE la vitesse : le temps de
travail du tour (.worktime) n'est pas une mesure de moteur, c'est ce que la
réponse a coûté. Une ligne qui ne portait que de la vitesse disparaît en
@@ -7399,16 +7402,32 @@ function renderStats(el, s){
if(!el||!s) return;
const parts=[];
const pt=s.prompt_tokens||s.prompt_tokens_total;
// Débit de lecture inconnu (API tierce sans timings) : on n'affiche que la taille.
if(pt) parts.push('prefill '+pt+' tok'+(s.prompt_per_second ? ' · '+s.prompt_per_second.toFixed(0)+' tok/s' : ''));
const pps=s.prompt_per_second ? ' · '+s.prompt_per_second.toFixed(0)+' tok/s' : '';
// Cache connu (cache_tokens présent, même à 0) : ce qui a été recalculé et ce
// qui venait du cache. Inconnu (moteur ancien, API tierce) : la taille seule,
// comme avant — un « 0 en cache » inventé serait pire que rien.
const cached=s.cache_tokens;
if(cached!=null && (s.prompt_tokens!=null || s.prompt_tokens_total)){
const fresh=s.prompt_tokens!=null ? s.prompt_tokens : Math.max(0,s.prompt_tokens_total-cached);
parts.push('prefill '+nfmt(fresh)+' nouveaux / '+nfmt(cached)+' en cache'+pps);
} else if(pt) parts.push('prefill '+pt+' tok'+pps); // débit inconnu (API tierce sans timings) : la taille seule
if(s.gen_tokens) parts.push('decode '+s.gen_tokens+' tok · '+(s.gen_per_second||0).toFixed(1)+' tok/s');
// Perte de cache signalée par le serveur (seuil relevé sur un modèle hybride,
// dont les points de reprise font perdre un peu à chaque étape).
if(s.lost_alert && s.lost) parts.push(nfmt(s.lost)+' recalculés'+(s.lost_after ? ' après '+s.lost_after : ''));
if(!parts.length) return;
// Réponse de l'assistant : ligne de mesures dédiée sous le texte (son étiquette
// est masquée dans cette mise en page). Bulle repliable : l'étiquette EST le
// bouton de repli, on y écrit comme avant.
if(el.classList.contains('collapsible')) setLabel(el, ['reasoning'].concat(workLabel()||[], parts).join(' · '));
else paintStats(el, parts.join(' · '));
else{
paintStats(el, parts.join(' · '));
const sl=el.querySelector(':scope > .statline');
if(sl) sl.classList.toggle('cache-miss', !!s.lost_alert);
}
}
// Milliers séparés par une espace fine insécable : « 41 230 », pas « 41230 ».
function nfmt(n){ return String(Math.round(n||0)).replace(/\B(?=(\d{3})+(?!\d))/g,'\u202f'); }
// --- Compteurs de bulle -----------------------------------------------------
// Ils sont PAR BULLE, jamais par tour. Un tour d'agent en ouvre une nouvelle
// après CHAQUE appel d'outil (le flux repasse T.contentEl à null) ; les
+19 -3
View File
@@ -237,16 +237,32 @@ function renderStats(el, s){
if(!el||!s) return;
const parts=[];
const pt=s.prompt_tokens||s.prompt_tokens_total;
// Débit de lecture inconnu (API tierce sans timings) : on n'affiche que la taille.
if(pt) parts.push('prefill '+pt+' tok'+(s.prompt_per_second ? ' · '+s.prompt_per_second.toFixed(0)+' tok/s' : ''));
const pps=s.prompt_per_second ? ' · '+s.prompt_per_second.toFixed(0)+' tok/s' : '';
// Cache connu (cache_tokens présent, même à 0) : ce qui a été recalculé et ce
// qui venait du cache. Inconnu (moteur ancien, API tierce) : la taille seule,
// comme avant — un « 0 en cache » inventé serait pire que rien.
const cached=s.cache_tokens;
if(cached!=null && (s.prompt_tokens!=null || s.prompt_tokens_total)){
const fresh=s.prompt_tokens!=null ? s.prompt_tokens : Math.max(0,s.prompt_tokens_total-cached);
parts.push('prefill '+nfmt(fresh)+' nouveaux / '+nfmt(cached)+' en cache'+pps);
} else if(pt) parts.push('prefill '+pt+' tok'+pps); // débit inconnu (API tierce sans timings) : la taille seule
if(s.gen_tokens) parts.push('decode '+s.gen_tokens+' tok · '+(s.gen_per_second||0).toFixed(1)+' tok/s');
// Perte de cache signalée par le serveur (seuil relevé sur un modèle hybride,
// dont les points de reprise font perdre un peu à chaque étape).
if(s.lost_alert && s.lost) parts.push(nfmt(s.lost)+' recalculés'+(s.lost_after ? ' après '+s.lost_after : ''));
if(!parts.length) return;
// Réponse de l'assistant : ligne de mesures dédiée sous le texte (son étiquette
// est masquée dans cette mise en page). Bulle repliable : l'étiquette EST le
// bouton de repli, on y écrit comme avant.
if(el.classList.contains('collapsible')) setLabel(el, ['reasoning'].concat(workLabel()||[], parts).join(' · '));
else paintStats(el, parts.join(' · '));
else{
paintStats(el, parts.join(' · '));
const sl=el.querySelector(':scope > .statline');
if(sl) sl.classList.toggle('cache-miss', !!s.lost_alert);
}
}
// Milliers séparés par une espace fine insécable : « 41 230 », pas « 41230 ».
function nfmt(n){ return String(Math.round(n||0)).replace(/\B(?=(\d{3})+(?!\d))/g,'\u202f'); }
// --- Compteurs de bulle -----------------------------------------------------
// Ils sont PAR BULLE, jamais par tour. Un tour d'agent en ouvre une nouvelle
// après CHAQUE appel d'outil (le flux repasse T.contentEl à null) ; les
+3
View File
@@ -1116,6 +1116,9 @@ h1{color:var(--text);font-weight:600;letter-spacing:-.01em}
#chat .msg .statline{margin-top:10px;font-family:var(--mono);font-size:10.5px;letter-spacing:.02em;
color:var(--dim);opacity:.75}
#chat .msg .statline:empty{display:none}
/* Cache de prompts perdu (gros recalcul) : la vitesse passe en ambre, la
ligne reste discrète mais ne se lit plus comme un tour ordinaire. */
#chat .msg .statline.cache-miss>.spd{color:var(--warn)}
/* « Masquer la vitesse de génération » ne masque QUE la vitesse : le temps de
travail du tour (.worktime) n'est pas une mesure de moteur, c'est ce que la
réponse a coûté. Une ligne qui ne portait que de la vitesse disparaît en
+1
View File
@@ -327,6 +327,7 @@ func newWebMux() *http.ServeMux {
api("/api/restart", svcHandler("restart"))
api("/api/bench", handleBench)
api("/api/bench/last", handleBenchLast)
api("/api/perf/summary", handlePerfSummary)
api("/api/chat", handleChat) // flux d'ABONNEMENT (SSE) : rejoue + suit le fil
api("/api/chat/send", handleChatSend) // envoie un message (lance la génération détachée)
api("/api/chat/upload", handleChatUpload) // dépose un fichier dans le workspace agent (joint au message suivant)