// 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 + OCR-Felder aus SRC-10). tableName ist über // tempTableNamePattern in Rebuild bereits geprüft, bevor diese Funktion // aufgerufen wird. MUSS bei jeder neuen Spalte in mail_documents // (migrations/000N_*.sql) mitgepflegt werden — sonst schlägt Reindex mit // "unknown column" fehl (siehe SRC-10, real so aufgetreten und behoben). 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, %s string attribute indexed, %s float)", tableName, FieldTenantSlug, FieldMessageID, FieldSubject, FieldBody, FieldAttachmentText, FieldSender, FieldMailbox, FieldAttachmentType, FieldTag, FieldSentAt, FieldOCRLanguage, FieldOCRConfidence, ) }