diff --git a/internal/loki/chat_compact.go b/internal/loki/chat_compact.go index 9e56f83..65b2c6c 100644 --- a/internal/loki/chat_compact.go +++ b/internal/loki/chat_compact.go @@ -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") } diff --git a/internal/loki/chat_conversation.go b/internal/loki/chat_conversation.go index 9b0c928..be4c705 100644 --- a/internal/loki/chat_conversation.go +++ b/internal/loki/chat_conversation.go @@ -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() diff --git a/internal/loki/code_subagent.go b/internal/loki/code_subagent.go index 3da21d8..0a18127 100644 --- a/internal/loki/code_subagent.go +++ b/internal/loki/code_subagent.go @@ -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() } diff --git a/internal/loki/code_verify.go b/internal/loki/code_verify.go index facb878..bbb4344 100644 --- a/internal/loki/code_verify.go +++ b/internal/loki/code_verify.go @@ -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 }) diff --git a/internal/loki/llm_bench.go b/internal/loki/llm_bench.go index 827b1b8..5b69ef4 100644 --- a/internal/loki/llm_bench.go +++ b/internal/loki/llm_bench.go @@ -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 diff --git a/internal/loki/llm_client.go b/internal/loki/llm_client.go index 698b490..ca2141d 100644 --- a/internal/loki/llm_client.go +++ b/internal/loki/llm_client.go @@ -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 } diff --git a/internal/loki/llm_oai.go b/internal/loki/llm_oai.go index 30df2d1..9737fee 100644 --- a/internal/loki/llm_oai.go +++ b/internal/loki/llm_oai.go @@ -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 } diff --git a/internal/loki/perf_log.go b/internal/loki/perf_log.go new file mode 100644 index 0000000..b0b701c --- /dev/null +++ b/internal/loki/perf_log.go @@ -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())) +} diff --git a/internal/loki/perf_log_test.go b/internal/loki/perf_log_test.go new file mode 100644 index 0000000..30848e6 --- /dev/null +++ b/internal/loki/perf_log_test.go @@ -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") + } +} diff --git a/internal/loki/tasks_run.go b/internal/loki/tasks_run.go index bda4494..bd9578c 100644 --- a/internal/loki/tasks_run.go +++ b/internal/loki/tasks_run.go @@ -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) } diff --git a/internal/loki/ui/index.html b/internal/loki/ui/index.html index ac1ce16..4a12130 100644 --- a/internal/loki/ui/index.html +++ b/internal/loki/ui/index.html @@ -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 diff --git a/internal/loki/ui/src/js/09-stream.js b/internal/loki/ui/src/js/09-stream.js index 0560e4b..ed115c4 100644 --- a/internal/loki/ui/src/js/09-stream.js +++ b/internal/loki/ui/src/js/09-stream.js @@ -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 diff --git a/internal/loki/ui/src/styles.css b/internal/loki/ui/src/styles.css index 44f568f..1a3071d 100644 --- a/internal/loki/ui/src/styles.css +++ b/internal/loki/ui/src/styles.css @@ -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 diff --git a/internal/loki/web_server.go b/internal/loki/web_server.go index bf8f67c..2611c36 100644 --- a/internal/loki/web_server.go +++ b/internal/loki/web_server.go @@ -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)