diff --git a/internal/loki/chat_persist.go b/internal/loki/chat_persist.go index 9846550..9ea3534 100644 --- a/internal/loki/chat_persist.go +++ b/internal/loki/chat_persist.go @@ -74,6 +74,11 @@ type convPersister struct { queued uint64 // tickets délivrés done uint64 // tickets traités (écrits, abandonnés ou en échec) running bool + + // Résultats d'outils du lot en cours d'écriture, et ceux d'entre eux qu'une + // suppression a annulés depuis : l'écrivain les saute dans sa transaction. + inflight map[string]bool + cancelled map[string]bool } var ( @@ -92,6 +97,9 @@ func newConvPersister() *convPersister { mem: map[string]string{}, written: map[string]uint64{}, deleted: map[string]bool{}, + + inflight: map[string]bool{}, + cancelled: map[string]bool{}, } p.cond = sync.NewCond(&p.mu) return p @@ -157,6 +165,7 @@ func (p *convPersister) run() { for _, j := range p.res { if !p.deleted[toolResConv(j.key)] { res = append(res, j) + p.inflight[j.key] = true } } p.snaps, p.res = map[string]*convSnap{}, nil @@ -166,6 +175,10 @@ func (p *convPersister) run() { wrote, err := persistCommitByPath(snaps, res) p.mu.Lock() + for _, j := range res { + delete(p.inflight, j.key) + delete(p.cancelled, j.key) + } for _, s := range wrote { if s.seq > p.written[s.id] { p.written[s.id] = s.seq @@ -225,6 +238,13 @@ func (p *convPersister) dropToolRes(prefix string) { } func (p *convPersister) dropToolResLocked(prefix string) { + // Déjà pris par l'écrivain : trop tard pour le retirer de la file, pas pour + // l'empêcher d'atteindre la base (voir toolResWanted). + for k := range p.inflight { + if strings.HasPrefix(k, prefix) { + p.cancelled[k] = true + } + } for k := range p.mem { if strings.HasPrefix(k, prefix) { delete(p.mem, k) @@ -239,6 +259,17 @@ func (p *convPersister) dropToolResLocked(prefix string) { p.res = kept } +// toolResWanted : ce résultat du lot en cours doit-il encore être écrit ? +// Appelé par l'écrivain DANS sa transaction. Une suppression annule avant +// d'ouvrir la sienne (deleteToolResultsFor) : soit l'écrivain voit +// l'annulation et saute le résultat, soit il l'a déjà écrit et la suppression, +// sérialisée après lui par bbolt, l'efface. Jamais de résultat qui renaît. +func (p *convPersister) toolResWanted(key string) bool { + p.mu.Lock() + defer p.mu.Unlock() + return !p.cancelled[key] +} + // toolResPending relit un résultat pas encore sur disque. func (p *convPersister) toolResPending(key string) (string, bool) { p.mu.Lock() @@ -338,6 +369,9 @@ func commitPersistBatch(path string, snaps []*convSnap, res []toolResJob) ([]*co return err } for _, j := range res { + if !persistQ.toolResWanted(j.key) { + continue // supprimé avec sa discussion pendant que le lot partait + } enc, err := encode([]byte(j.plain)) if err != nil { return err diff --git a/internal/loki/chat_persist_test.go b/internal/loki/chat_persist_test.go index e3bc843..cc6d3fb 100644 --- a/internal/loki/chat_persist_test.go +++ b/internal/loki/chat_persist_test.go @@ -187,6 +187,29 @@ func TestToolResultSupprimeAvantEcriture(t *testing.T) { } } +// Un résultat que l'écrivain a DÉJÀ pris dans son lot, supprimé avant que ce +// lot n'atteigne la base, ne renaît pas après la suppression. C'était la cause +// des échecs intermittents de TestToolResultsSaveLoadDelete : la suppression ne +// retirait que la file d'attente, et l'écriture en vol passait après elle. +func TestToolResultSupprimePendantEcriture(t *testing.T) { + testHome(t) + entered, release := gatePersist(t) + id := saveToolResult(strings.Repeat("a", 3000)) + waitEntered(t, entered) // le lot est pris, pas encore écrit + deleteToolResultsFor(toolResConv(id)) + close(release) + persistQ.flush() + if _, ok := loadToolResult(id); ok { + t.Fatal("un résultat supprimé pendant son écriture a été écrit quand même") + } + // L'annulation ne vaut que pour ce lot : un nouveau résultat s'écrit. + again := saveToolResult(strings.Repeat("b", 3000)) + persistQ.flush() + if _, ok := loadToolResult(again); !ok { + t.Fatal("annulation restée collée : le résultat suivant n'a pas été écrit") + } +} + // compactLog rend exactement ce que donne la compaction de fin de tour, sans // toucher au journal d'origine. func TestCompactLogPurEgalFinDeTour(t *testing.T) { diff --git a/internal/loki/tool_results.go b/internal/loki/tool_results.go index 391f416..7f13915 100644 --- a/internal/loki/tool_results.go +++ b/internal/loki/tool_results.go @@ -89,6 +89,10 @@ func deleteToolResultsFor(sid string) { return } prefix := sid + "." + // D'abord l'écrivain (chat_persist.go) : ce qui l'attend est jeté, et ce + // qu'il a déjà pris dans son lot est annulé — sans ça, il l'écrirait après + // l'effacement ci-dessous et le résultat renaîtrait. Annulé AVANT notre + // transaction : l'écrivain vérifie dans la sienne, et bbolt les sérialise. persistQ.dropToolRes(prefix) _ = update(bkToolRes, func(b *bolt.Bucket) error { c := b.Cursor()