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
327 lines
11 KiB
Go
327 lines
11 KiB
Go
package search
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
_ "embed"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
// 0001_mail_documents.sql: statisches, versioniertes Schema
|
|
// (Akzeptanzkriterium 1). Feldnamen hier UND in fields.go müssen
|
|
// deckungsgleich bleiben — die Konstanten in fields.go sind die einzige
|
|
// Stelle, aus der Go-Code Feldnamen für Schreib-/Lesezugriffe bezieht.
|
|
// Manticores SQL-Parser unterstützt keine "--"-Kommentare, daher bleibt
|
|
// die eingebettete Datei selbst kommentarfrei.
|
|
//
|
|
//go:embed migrations/0001_mail_documents.sql
|
|
var schemaMigration string
|
|
|
|
// SRC-05: Facettenfelder als eigene, nummerierte ALTER-Migrationen
|
|
// nachgezogen (Manticore erlaubt nur eine Spalte je ALTER TABLE ADD
|
|
// COLUMN-Anweisung). Reihenfolge ist die Anwendungsreihenfolge.
|
|
//
|
|
//go:embed migrations/0002_mail_documents_facets.sql
|
|
var migrationAddSender string
|
|
|
|
//go:embed migrations/0003_mail_documents_mailbox.sql
|
|
var migrationAddMailbox string
|
|
|
|
//go:embed migrations/0004_mail_documents_attachment_type.sql
|
|
var migrationAddAttachmentType string
|
|
|
|
//go:embed migrations/0005_mail_documents_tag.sql
|
|
var migrationAddTag string
|
|
|
|
var facetMigrations = []string{
|
|
migrationAddSender,
|
|
migrationAddMailbox,
|
|
migrationAddAttachmentType,
|
|
migrationAddTag,
|
|
}
|
|
|
|
// Client spricht ausschließlich über die strukturierte Manticore-HTTP-
|
|
// JSON-API (kein String-Zusammenbau von SQL-Klauseln, siehe fields.go).
|
|
// Die SQL-Schnittstelle wird nur für EnsureSchema verwendet, und dort
|
|
// ausschließlich mit dem statischen, eingebetteten Migrationstext —
|
|
// niemals mit zur Laufzeit zusammengesetzten Werten.
|
|
type Client struct {
|
|
baseURL string
|
|
http *http.Client
|
|
}
|
|
|
|
func NewClient(baseURL string) *Client {
|
|
return &Client{
|
|
baseURL: strings.TrimRight(baseURL, "/"),
|
|
http: &http.Client{Timeout: 10 * time.Second},
|
|
}
|
|
}
|
|
|
|
// EnsureSchema legt den Index gemäß dem versionierten, statischen
|
|
// Migrationstext an (Akzeptanzkriterium 1) und zieht die Facettenfelder
|
|
// (SRC-05) idempotent nach.
|
|
func (c *Client) EnsureSchema(ctx context.Context) error {
|
|
if err := c.runSchemaSQL(ctx, schemaMigration); err != nil {
|
|
return fmt.Errorf("search: schema anlegen: %w", err)
|
|
}
|
|
for _, migration := range facetMigrations {
|
|
if err := c.runSchemaSQL(ctx, migration); err != nil {
|
|
// Manticore meldet bei erneutem ADD COLUMN "field already in
|
|
// schema" — kein Fehler, sondern der bereits angewendete
|
|
// Migrationsschritt (Idempotenz, gleiche CREATE-TABLE-IF-NOT-
|
|
// EXISTS-Konvention wie das Basisschema).
|
|
if strings.Contains(err.Error(), "already in schema") {
|
|
continue
|
|
}
|
|
return fmt.Errorf("search: facettenfeld-migration: %w", err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (c *Client) runSchemaSQL(ctx context.Context, query string) error {
|
|
form := "query=" + query
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/sql?mode=raw", strings.NewReader(form))
|
|
if err != nil {
|
|
return 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 fmt.Errorf("ausführen: %w", err)
|
|
}
|
|
defer func() { _ = resp.Body.Close() }()
|
|
body, _ := io.ReadAll(resp.Body)
|
|
if resp.StatusCode != http.StatusOK {
|
|
return fmt.Errorf("status %d: %s", resp.StatusCode, string(body))
|
|
}
|
|
if strings.Contains(string(body), `"error":"`) && !strings.Contains(string(body), `"error":""`) {
|
|
return fmt.Errorf("%s", string(body))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Document ist ein Mail-Suchdokument. Feldnamen im JSON-Tag entsprechen
|
|
// exakt den Konstanten in fields.go.
|
|
type Document struct {
|
|
ID uint64 `json:"-"`
|
|
TenantSlug string `json:"tenant_slug"`
|
|
MessageID string `json:"message_id"`
|
|
Subject string `json:"subject"`
|
|
Body string `json:"body"`
|
|
AttachmentText string `json:"attachment_text"`
|
|
SentAtUnixEpoch int64 `json:"sent_at"`
|
|
// Facettenfelder (SRC-05), optional — leerer String bedeutet "kein Wert".
|
|
Sender string `json:"sender"`
|
|
Mailbox string `json:"mailbox"`
|
|
AttachmentType string `json:"attachment_type"`
|
|
Tag string `json:"tag"`
|
|
}
|
|
|
|
// Index legt/ersetzt ein Suchdokument (Akzeptanzkriterium 2: Schreibzugriff
|
|
// ausschließlich über statische, vordefinierte Feldnamen aus dem
|
|
// Document-Struct — kein dynamischer Feldname möglich).
|
|
func (c *Client) Index(ctx context.Context, doc Document) error {
|
|
payload := map[string]any{
|
|
"index": IndexName,
|
|
"id": doc.ID,
|
|
"doc": doc,
|
|
}
|
|
body, err := json.Marshal(payload)
|
|
if err != nil {
|
|
return fmt.Errorf("search: dokument serialisieren: %w", err)
|
|
}
|
|
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/replace", bytes.NewReader(body))
|
|
if err != nil {
|
|
return fmt.Errorf("search: index-anfrage bauen: %w", err)
|
|
}
|
|
req.Header.Set("Content-Type", "application/json")
|
|
|
|
resp, err := c.http.Do(req)
|
|
if err != nil {
|
|
return fmt.Errorf("search: dokument indexieren: %w", err)
|
|
}
|
|
defer func() { _ = resp.Body.Close() }()
|
|
respBody, _ := io.ReadAll(resp.Body)
|
|
if resp.StatusCode != http.StatusOK {
|
|
return fmt.Errorf("search: dokument indexieren, status %d: %s", resp.StatusCode, string(respBody))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Delete entfernt ein Suchdokument anhand seiner ID (SRC-02
|
|
// Akzeptanzkriterium 2: Löschungen werden im Index nachgezogen). Löschen
|
|
// eines nicht (mehr) vorhandenen Dokuments ist kein Fehler (idempotent).
|
|
func (c *Client) Delete(ctx context.Context, id uint64) error {
|
|
payload := map[string]any{
|
|
"index": IndexName,
|
|
"id": id,
|
|
}
|
|
body, err := json.Marshal(payload)
|
|
if err != nil {
|
|
return fmt.Errorf("search: lösch-anfrage serialisieren: %w", err)
|
|
}
|
|
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/delete", bytes.NewReader(body))
|
|
if err != nil {
|
|
return fmt.Errorf("search: lösch-anfrage bauen: %w", err)
|
|
}
|
|
req.Header.Set("Content-Type", "application/json")
|
|
|
|
resp, err := c.http.Do(req)
|
|
if err != nil {
|
|
return fmt.Errorf("search: dokument löschen: %w", err)
|
|
}
|
|
defer func() { _ = resp.Body.Close() }()
|
|
respBody, _ := io.ReadAll(resp.Body)
|
|
if resp.StatusCode != http.StatusOK {
|
|
return fmt.Errorf("search: dokument löschen, status %d: %s", resp.StatusCode, string(respBody))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Result ist ein Suchtreffer.
|
|
type Result struct {
|
|
MessageID string
|
|
Subject string
|
|
Score int64
|
|
SentAtUnixEpoch int64
|
|
}
|
|
|
|
// fieldWeights gewichtet Betreff höher als Text, Anhangstext am
|
|
// niedrigsten (SRC-03 Akzeptanzkriterium 1: Ranking berücksichtigt u.a.
|
|
// Anhangstreffer) — statische Konstanten, keine dynamischen Feldnamen.
|
|
var fieldWeights = map[string]any{
|
|
FieldSubject: 10,
|
|
FieldBody: 3,
|
|
FieldAttachmentText: 1,
|
|
}
|
|
|
|
// doSearchWithSwapRetry führt eine /search-Anfrage aus und wiederholt sie
|
|
// bis zu zweimal mit kurzer Pause, falls Manticore "unknown local table"
|
|
// meldet (SRC-09 Akzeptanzkriterium 3: der Reindex-Umschaltmoment
|
|
// RENAME-alte-Tabelle-weg/RENAME-neue-Tabelle-rein hat ein extrem kurzes
|
|
// Zeitfenster ohne existierende mail_documents-Tabelle — dieser Retry
|
|
// überbrückt es, statt eine Suchanfrage in genau diesem Moment fehlschlagen
|
|
// zu lassen).
|
|
func (c *Client) doSearchWithSwapRetry(ctx context.Context, body []byte) ([]byte, error) {
|
|
const maxAttempts = 3
|
|
var lastErr error
|
|
for attempt := 0; attempt < maxAttempts; attempt++ {
|
|
if attempt > 0 {
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil, ctx.Err()
|
|
case <-time.After(20 * time.Millisecond):
|
|
}
|
|
}
|
|
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/search", bytes.NewReader(body))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("search: suchanfrage bauen: %w", err)
|
|
}
|
|
req.Header.Set("Content-Type", "application/json")
|
|
|
|
resp, err := c.http.Do(req)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("search: suche ausführen: %w", err)
|
|
}
|
|
respBody, readErr := io.ReadAll(resp.Body)
|
|
_ = resp.Body.Close()
|
|
if readErr != nil {
|
|
return nil, fmt.Errorf("search: antwort lesen: %w", readErr)
|
|
}
|
|
if strings.Contains(string(respBody), "unknown local table") {
|
|
lastErr = fmt.Errorf("search: suche, status %d: %s", resp.StatusCode, string(respBody))
|
|
continue
|
|
}
|
|
if resp.StatusCode != http.StatusOK {
|
|
return nil, fmt.Errorf("search: suche, status %d: %s", resp.StatusCode, string(respBody))
|
|
}
|
|
return respBody, nil
|
|
}
|
|
return nil, lastErr
|
|
}
|
|
|
|
// Search sucht queryText innerhalb der Volltextfelder, strikt begrenzt auf
|
|
// den Mandanten tenantSlug (Akzeptanzkriterium 2: mandantengetrennt
|
|
// abfragbar) — der Tenant-Filter läuft über ein strukturiertes "equals"-
|
|
// Feld der JSON-API, niemals über eine interpolierte WHERE-Klausel.
|
|
//
|
|
// queryText nutzt Manticores erweiterte Abfragesyntax über den
|
|
// query_string-Klausel-Typ (Akzeptanzkriterium 3: Phrasensuche mit
|
|
// Anführungszeichen, Ausschluss mit vorangestelltem "-") — der Wert landet
|
|
// als reiner JSON-String-Wert, es gibt dabei keinerlei dynamischen
|
|
// Feld-/Tabellennamen, der beeinflusst werden könnte. Ergebnisse kommen
|
|
// von Manticore bereits nach Relevanz (BM25, gewichtet über fieldWeights)
|
|
// absteigend sortiert zurück (Akzeptanzkriterium 1).
|
|
const searchResultLimit = 1000
|
|
|
|
func (c *Client) Search(ctx context.Context, tenantSlug, queryText string) ([]Result, error) {
|
|
payload := map[string]any{
|
|
"index": IndexName,
|
|
"query": map[string]any{
|
|
"bool": map[string]any{
|
|
"must": []map[string]any{
|
|
{"equals": map[string]any{FieldTenantSlug: tenantSlug}},
|
|
{"query_string": queryText},
|
|
},
|
|
},
|
|
},
|
|
"options": map[string]any{
|
|
"field_weights": fieldWeights,
|
|
},
|
|
// Ohne explizites limit begrenzt Manticore standardmäßig auf 20
|
|
// Treffer — bei Testkorpora bis 1000 Dokumenten (SRC-03) blieb das
|
|
// bisher unbemerkt, da nur auf das Vorhandensein einzelner Treffer
|
|
// geprüft wurde, nicht auf die Gesamtzahl. searchResultLimit deckt
|
|
// realistische Trefferlisten ab, ohne unbegrenzt zu sein.
|
|
"limit": searchResultLimit,
|
|
}
|
|
body, err := json.Marshal(payload)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("search: suchanfrage serialisieren: %w", err)
|
|
}
|
|
|
|
respBody, err := c.doSearchWithSwapRetry(ctx, body)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
var parsed searchResponse
|
|
if err := json.Unmarshal(respBody, &parsed); err != nil {
|
|
return nil, fmt.Errorf("search: antwort parsen: %w", err)
|
|
}
|
|
|
|
results := make([]Result, 0, len(parsed.Hits.Hits))
|
|
for _, hit := range parsed.Hits.Hits {
|
|
results = append(results, Result{
|
|
MessageID: hit.Source.MessageID,
|
|
Subject: hit.Source.Subject,
|
|
Score: hit.Score,
|
|
SentAtUnixEpoch: hit.Source.SentAtUnixEpoch,
|
|
})
|
|
}
|
|
return results, nil
|
|
}
|
|
|
|
type searchResponse struct {
|
|
Hits struct {
|
|
Hits []struct {
|
|
Score int64 `json:"_score"`
|
|
Source struct {
|
|
MessageID string `json:"message_id"`
|
|
Subject string `json:"subject"`
|
|
SentAtUnixEpoch int64 `json:"sent_at"`
|
|
} `json:"_source"`
|
|
} `json:"hits"`
|
|
} `json:"hits"`
|
|
}
|