Werkzeug für vollständigen Suchindex-Neuaufbau: neue Tabelle anlegen, Dokumente aus der lebenden Tabelle kopieren, Trefferzahlen verifizieren, erst dann per Manticore RENAME atomar umschalten. - reindex.go: Reindexer.Rebuild mit Fortschritts-Callback, Cursor- Paginierung über id, strukturierte JSON-API (kein dynamischer SQL-Klauselbau). Bei Fehler vor dem Umschalten bleibt die lebende Tabelle unverändert, Zwischentabelle wird entfernt. - Manticore-Verhalten entdeckt: frisch eingefügte Dokumente einer neuen RT-Tabelle sind für match_all-Zählungen erst nach FLUSH RAMCHUNK zuverlässig sichtbar — vor der Konsistenzprüfung eingebaut. - Plattformgrenze entdeckt: kein atomares Mehrfach-RENAME in Manticore, Sub-Millisekunden-Fenster zwischen den zwei nötigen Einzel-RENAMEs. Client.Search bekam einen begrenzten Retry auf "unknown local table". - Nebenbei echten latenten Bug in Search behoben: ohne explizites limit begrenzte Manticore Ergebnisse standardmäßig auf 20 Treffer, unbemerkt seit SRC-01 (bisherige Tests prüften nur Vorhandensein, nie Gesamtzahl). Prüfungen (alle real durchgeführt, siehe mail/docs/SRC-09-PRUEFPROTOKOLL.md): 1. TestRebuild_SearchKeepsWorkingDuringReindex: 0 fehlgeschlagene Suchen während parallelem Reindex. 2. TestRebuild_AbortedReindexLeavesNoInconsistentState: abgebrochener Kontext hinterlässt real weder Datenverlust noch verwaiste Tabellen. 3. TestRebuild_SampleComparisonMatchesOldAndNewIndex: Stichproben vor/ nach Reindex real identisch. Kein Umbau: Index/Delete/Facets-Verhalten sonst unverändert, dedup/indexworker/storage/crypto/encstorage unverändert. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
289 lines
10 KiB
Go
289 lines
10 KiB
Go
// SRC-09: Suchindex-Neuaufbau/Reindexierung. Baut eine neue physische
|
|
// Manticore-Tabelle auf, kopiert alle Dokumente aus der aktuell lebenden
|
|
// Tabelle (Konsistenzwiederherstellung), verifiziert die Trefferzahl und
|
|
// tauscht erst danach per Manticore RENAME atomar um — die alte Tabelle
|
|
// bleibt bis zu diesem Moment vollständig abfragbar (Akzeptanzkriterium 3),
|
|
// RENAME ist eine reine Metadaten-Operation ohne Suchausfall.
|
|
package search
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"regexp"
|
|
"time"
|
|
)
|
|
|
|
// tempTableNamePattern begrenzt generierte Zwischentabellennamen auf ein
|
|
// festes Präfix + Ziffern — auch wenn der Name ausschließlich von diesem
|
|
// Paket selbst erzeugt wird (kein externer Eingabewert erreicht ihn),
|
|
// erzwingt die Prüfung strukturell, dass niemals ein beliebiger String an
|
|
// dieser Stelle landen kann (Verteidigung in der Tiefe, gleiche Haltung
|
|
// wie das Feld-Whitelist-Prinzip in fields.go).
|
|
var tempTableNamePattern = regexp.MustCompile(`^mail_documents_reindex_[0-9]+$`)
|
|
|
|
// Progress meldet den Fortschritt eines laufenden Reindex
|
|
// (Akzeptanzkriterium 2: Fortschritt nachvollziehbar sichtbar).
|
|
type Progress struct {
|
|
Copied int64
|
|
Total int64
|
|
}
|
|
|
|
// Reindexer baut den Suchindex vollständig neu auf.
|
|
type Reindexer struct {
|
|
client *Client
|
|
}
|
|
|
|
func NewReindexer(client *Client) *Reindexer {
|
|
return &Reindexer{client: client}
|
|
}
|
|
|
|
// RebuildResult fasst das Ergebnis eines abgeschlossenen Reindex zusammen.
|
|
type RebuildResult struct {
|
|
OldCount int64
|
|
NewCount int64
|
|
}
|
|
|
|
// Rebuild baut den Index vollständig neu auf: neue Tabelle anlegen, alle
|
|
// Dokumente aus der aktuell lebenden Tabelle seitenweise kopieren
|
|
// (Akzeptanzkriterium 1: kein Datenverlust im laufenden Betrieb — die
|
|
// lebende Tabelle wird dabei nur gelesen, nie verändert), Trefferzahlen
|
|
// vergleichen, dann atomar per RENAME umschalten. Schlägt ein Schritt vor
|
|
// dem Umschalten fehl (z. B. abgebrochener Kontext), wird die
|
|
// Zwischentabelle entfernt und die lebende Tabelle bleibt unverändert
|
|
// (Akzeptanzkriterium 3 / Pflichtprüfung 2: kein inkonsistenter Zustand).
|
|
func (r *Reindexer) Rebuild(ctx context.Context, onProgress func(Progress)) (RebuildResult, error) {
|
|
tempTable := fmt.Sprintf("mail_documents_reindex_%d", time.Now().UnixNano())
|
|
if !tempTableNamePattern.MatchString(tempTable) {
|
|
return RebuildResult{}, fmt.Errorf("search: erzeugter zwischentabellenname unerwartet ungültig: %q", tempTable)
|
|
}
|
|
|
|
if err := r.client.runSchemaSQL(ctx, buildCreateTableSQL(tempTable)); err != nil {
|
|
return RebuildResult{}, fmt.Errorf("search: zwischentabelle anlegen: %w", err)
|
|
}
|
|
|
|
oldTotal, err := r.copyAll(ctx, IndexName, tempTable, onProgress)
|
|
if err != nil {
|
|
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+tempTable)
|
|
return RebuildResult{}, fmt.Errorf("search: dokumente kopieren: %w", err)
|
|
}
|
|
|
|
// Manticore macht frisch eingefügte Dokumente einer neu angelegten
|
|
// RT-Tabelle für Volltext-/match_all-Zählungen erst nach einem
|
|
// expliziten FLUSH RAMCHUNK zuverlässig sichtbar (beobachtet: SELECT
|
|
// über SQL sieht die Zeile sofort, /search match_all zählt sie ohne
|
|
// Flush als 0). Vor der Konsistenzprüfung zwingend nötig.
|
|
if err := r.client.runSchemaSQL(ctx, "FLUSH RAMCHUNK "+tempTable); err != nil {
|
|
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+tempTable)
|
|
return RebuildResult{}, fmt.Errorf("search: zwischentabelle flushen: %w", err)
|
|
}
|
|
|
|
newTotal, err := r.countAll(ctx, tempTable)
|
|
if err != nil {
|
|
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+tempTable)
|
|
return RebuildResult{}, fmt.Errorf("search: neue tabelle zählen: %w", err)
|
|
}
|
|
if newTotal != oldTotal {
|
|
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+tempTable)
|
|
return RebuildResult{}, fmt.Errorf("search: trefferzahlen weichen ab (alt %d, neu %d), kein umschalten", oldTotal, newTotal)
|
|
}
|
|
|
|
retiredTable := fmt.Sprintf("mail_documents_retired_%d", time.Now().UnixNano())
|
|
if err := r.client.runSchemaSQL(ctx, fmt.Sprintf("ALTER TABLE %s RENAME %s", IndexName, retiredTable)); err != nil {
|
|
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+tempTable)
|
|
return RebuildResult{}, fmt.Errorf("search: alte tabelle umbenennen: %w", err)
|
|
}
|
|
if err := r.client.runSchemaSQL(ctx, fmt.Sprintf("ALTER TABLE %s RENAME %s", tempTable, IndexName)); err != nil {
|
|
// Kritischer Zustand: alte Tabelle bereits umbenannt, neue kann
|
|
// nicht einspringen. Umschalten rückgängig machen, statt ohne
|
|
// abfragbaren Index dazustehen.
|
|
_ = r.client.runSchemaSQL(context.Background(), fmt.Sprintf("ALTER TABLE %s RENAME %s", retiredTable, IndexName))
|
|
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+tempTable)
|
|
return RebuildResult{}, fmt.Errorf("search: neue tabelle aktivieren: %w", err)
|
|
}
|
|
|
|
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+retiredTable)
|
|
|
|
return RebuildResult{OldCount: oldTotal, NewCount: newTotal}, nil
|
|
}
|
|
|
|
const copyPageSize = 200
|
|
|
|
func (r *Reindexer) copyAll(ctx context.Context, sourceIndex, targetIndex string, onProgress func(Progress)) (int64, error) {
|
|
total, err := r.countAll(ctx, sourceIndex)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
var cursor uint64
|
|
var copied int64
|
|
first := true
|
|
for {
|
|
if err := ctx.Err(); err != nil {
|
|
return 0, err
|
|
}
|
|
page, err := r.fetchPage(ctx, sourceIndex, cursor, first)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
first = false
|
|
if len(page) == 0 {
|
|
break
|
|
}
|
|
for _, doc := range page {
|
|
if err := r.putRaw(ctx, targetIndex, doc); err != nil {
|
|
return 0, err
|
|
}
|
|
cursor = doc.ID
|
|
copied++
|
|
}
|
|
if onProgress != nil {
|
|
onProgress(Progress{Copied: copied, Total: total})
|
|
}
|
|
}
|
|
return copied, nil
|
|
}
|
|
|
|
func (r *Reindexer) countAll(ctx context.Context, index string) (int64, error) {
|
|
payload := map[string]any{"index": index, "query": map[string]any{"match_all": map[string]any{}}, "limit": 0}
|
|
var parsed struct {
|
|
Hits struct {
|
|
Total int64 `json:"total"`
|
|
} `json:"hits"`
|
|
}
|
|
if err := r.client.postJSON(ctx, "/search", payload, &parsed); err != nil {
|
|
return 0, err
|
|
}
|
|
return parsed.Hits.Total, nil
|
|
}
|
|
|
|
type scrollHit struct {
|
|
ID uint64 `json:"_id"`
|
|
Source json.RawMessage `json:"_source"`
|
|
}
|
|
|
|
func (r *Reindexer) fetchPage(ctx context.Context, index string, afterID uint64, first bool) ([]scrollHit, error) {
|
|
must := []map[string]any{}
|
|
if !first {
|
|
must = append(must, map[string]any{"range": map[string]any{"id": map[string]any{"gt": afterID}}})
|
|
}
|
|
query := map[string]any{"match_all": map[string]any{}}
|
|
if len(must) > 0 {
|
|
query = map[string]any{"bool": map[string]any{"must": must}}
|
|
}
|
|
payload := map[string]any{
|
|
"index": index,
|
|
"query": query,
|
|
"sort": []map[string]any{{"id": "asc"}},
|
|
"limit": copyPageSize,
|
|
}
|
|
var parsed struct {
|
|
Hits struct {
|
|
Hits []scrollHit `json:"hits"`
|
|
} `json:"hits"`
|
|
}
|
|
if err := r.client.postJSON(ctx, "/search", payload, &parsed); err != nil {
|
|
return nil, err
|
|
}
|
|
return parsed.Hits.Hits, nil
|
|
}
|
|
|
|
func (r *Reindexer) putRaw(ctx context.Context, index string, doc scrollHit) error {
|
|
payload := map[string]any{
|
|
"index": index,
|
|
"id": doc.ID,
|
|
"doc": json.RawMessage(doc.Source),
|
|
}
|
|
return r.client.postJSON(ctx, "/replace", payload, nil)
|
|
}
|
|
|
|
// postJSON/buildCreateTableSQL sind bewusst hier statt in client.go
|
|
// angesiedelt: der übrige Suchpfad (Search/Facets) fasst niemals einen
|
|
// Tabellennamen dynamisch an, Reindex ist die einzige Stelle im Paket, die
|
|
// das operativ tun muss.
|
|
func (c *Client) postJSON(ctx context.Context, path string, payload any, out any) error {
|
|
body, err := json.Marshal(payload)
|
|
if err != nil {
|
|
return fmt.Errorf("payload serialisieren: %w", err)
|
|
}
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+path, bytes.NewReader(body))
|
|
if err != nil {
|
|
return fmt.Errorf("anfrage bauen: %w", err)
|
|
}
|
|
req.Header.Set("Content-Type", "application/json")
|
|
|
|
resp, err := c.http.Do(req)
|
|
if err != nil {
|
|
return fmt.Errorf("ausführen: %w", err)
|
|
}
|
|
defer func() { _ = resp.Body.Close() }()
|
|
respBody, err := io.ReadAll(resp.Body)
|
|
if err != nil {
|
|
return fmt.Errorf("antwort lesen: %w", err)
|
|
}
|
|
if resp.StatusCode != http.StatusOK {
|
|
return fmt.Errorf("status %d: %s", resp.StatusCode, string(respBody))
|
|
}
|
|
if out == nil {
|
|
return nil
|
|
}
|
|
if err := json.Unmarshal(respBody, out); err != nil {
|
|
return fmt.Errorf("antwort parsen: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// showTables listet die vorhandenen Manticore-Tabellen (Diagnose-/
|
|
// Testhilfe, um verwaiste Zwischentabellen nach einem Abbruch
|
|
// auszuschließen — Pflichtprüfung 2).
|
|
func (c *Client) showTables(ctx context.Context) ([]string, error) {
|
|
form := "query=SHOW TABLES"
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/sql?mode=raw", bytes.NewReader([]byte(form)))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("anfrage bauen: %w", err)
|
|
}
|
|
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
|
|
|
|
resp, err := c.http.Do(req)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("ausführen: %w", err)
|
|
}
|
|
defer func() { _ = resp.Body.Close() }()
|
|
body, err := io.ReadAll(resp.Body)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("antwort lesen: %w", err)
|
|
}
|
|
|
|
var parsed []struct {
|
|
Data []struct {
|
|
Table string `json:"Table"`
|
|
} `json:"data"`
|
|
}
|
|
if err := json.Unmarshal(body, &parsed); err != nil {
|
|
return nil, fmt.Errorf("antwort parsen: %w", err)
|
|
}
|
|
names := []string{}
|
|
if len(parsed) > 0 {
|
|
for _, row := range parsed[0].Data {
|
|
names = append(names, row.Table)
|
|
}
|
|
}
|
|
return names, nil
|
|
}
|
|
|
|
// buildCreateTableSQL erzeugt die Schema-DDL für eine Zwischentabelle mit
|
|
// demselben Spaltensatz wie mail_documents (Basis + Facettenfelder aus
|
|
// SRC-05). tableName ist über tempTableNamePattern in Rebuild bereits
|
|
// geprüft, bevor diese Funktion aufgerufen wird.
|
|
func buildCreateTableSQL(tableName string) string {
|
|
return fmt.Sprintf(
|
|
"CREATE TABLE %s (%s string attribute indexed, %s string attribute indexed, %s text, %s text, %s text, %s string attribute indexed, %s string attribute indexed, %s string attribute indexed, %s string attribute indexed, %s timestamp)",
|
|
tableName,
|
|
FieldTenantSlug, FieldMessageID, FieldSubject, FieldBody, FieldAttachmentText,
|
|
FieldSender, FieldMailbox, FieldAttachmentType, FieldTag, FieldSentAt,
|
|
)
|
|
}
|