* feat(metadata): improve match accuracy and localized titles * fix(metadata): address matching review findings * test(catalog): align empty alias snapshot scope --------- Co-authored-by: Quick104 <31828688+Quick104@users.noreply.github.com>
990 lines
34 KiB
Go
990 lines
34 KiB
Go
package catalog
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"log/slog"
|
|
"math"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/Silo-Server/silo-server/internal/embeddingvectors"
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
"github.com/pgvector/pgvector-go"
|
|
)
|
|
|
|
type SearchIndexProgressReporter interface {
|
|
Report(percent float64, message string)
|
|
SetResultData(data json.RawMessage)
|
|
}
|
|
|
|
type CatalogSearchIndexSyncStats struct {
|
|
Configured bool `json:"configured"`
|
|
Skipped bool `json:"skipped"`
|
|
Reason string `json:"reason,omitempty"`
|
|
Events int `json:"events"`
|
|
Upserted int `json:"upserted"`
|
|
Deleted int `json:"deleted"`
|
|
ActiveIndexUID string `json:"active_index_uid,omitempty"`
|
|
DocumentCount int `json:"document_count"`
|
|
VectorDocCount int `json:"vector_document_count"`
|
|
LastProcessedID int64 `json:"last_processed_event_id,omitempty"`
|
|
}
|
|
|
|
type CatalogSearchIndexRebuildStats struct {
|
|
Configured bool `json:"configured"`
|
|
Skipped bool `json:"skipped"`
|
|
Reason string `json:"reason,omitempty"`
|
|
ActiveIndexUID string `json:"active_index_uid,omitempty"`
|
|
DocumentCount int `json:"document_count"`
|
|
VectorDocCount int `json:"vector_document_count"`
|
|
RemovedIndexes int `json:"removed_indexes"`
|
|
}
|
|
|
|
type queuedMeilisearchTask struct {
|
|
task meilisearchTaskRef
|
|
docCount int
|
|
vecCount int
|
|
}
|
|
|
|
const meilisearchMaxDocumentPayloadBytes = 80 * 1024 * 1024
|
|
|
|
// meilisearchIndexingTimeout bounds a single indexing HTTP call (document
|
|
// batches run up to meilisearchMaxDocumentPayloadBytes). It is deliberately
|
|
// independent of catalog.search.meilisearch.timeout_ms: that setting protects
|
|
// the interactive search hot path and defaults to 800ms, which is far too
|
|
// tight to upload a multi-megabyte rebuild batch to a non-loopback
|
|
// Meilisearch — and raising it to make indexing work would loosen search
|
|
// fallback latency at the same time.
|
|
const meilisearchIndexingTimeout = 2 * time.Minute
|
|
|
|
// catalogSearchExcludeMangaChaptersSQL excludes per-chapter manga rows from
|
|
// catalog search documents and semantic coverage counts; chapters are reached
|
|
// through their parent series and would otherwise flood the index. The
|
|
// predicate expects media_items to be aliased as `mi`.
|
|
const catalogSearchExcludeMangaChaptersSQL = `NOT EXISTS (SELECT 1 FROM manga_chapters mc WHERE mc.chapter_content_id = mi.content_id)`
|
|
|
|
type CatalogSearchIndexer struct {
|
|
pool *pgxpool.Pool
|
|
settingsStore SettingsStore
|
|
events *SearchIndexEventRepository
|
|
}
|
|
|
|
func NewCatalogSearchIndexer(pool *pgxpool.Pool, settingsStore SettingsStore) *CatalogSearchIndexer {
|
|
return &CatalogSearchIndexer{
|
|
pool: pool,
|
|
settingsStore: settingsStore,
|
|
events: NewSearchIndexEventRepository(pool),
|
|
}
|
|
}
|
|
|
|
func (i *CatalogSearchIndexer) ShouldSyncRun(ctx context.Context) (bool, error) {
|
|
settings, ok, err := i.loadMeilisearchRuntime(ctx)
|
|
if err != nil || !ok || settings.Provider != SearchProviderMeilisearch {
|
|
return false, err
|
|
}
|
|
state, err := i.events.GetState(ctx, SearchProviderMeilisearch)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
if state.ActiveIndexUID == "" || state.SchemaVersion != catalogSearchMeilisearchSchemaVersion(settings.Embedder, settings.IndexTypes, settings.SemanticEnabled, settings.BinaryQuantized) {
|
|
return false, nil
|
|
}
|
|
pending, err := i.events.PendingCount(ctx, SearchProviderMeilisearch)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
return pending > 0, nil
|
|
}
|
|
|
|
func (i *CatalogSearchIndexer) SyncOutbox(ctx context.Context, progress SearchIndexProgressReporter) (CatalogSearchIndexSyncStats, error) {
|
|
stats := CatalogSearchIndexSyncStats{}
|
|
settings, client, ok, err := i.loadClient(ctx)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
if !ok {
|
|
stats.Skipped = true
|
|
stats.Reason = "meilisearch is not configured"
|
|
setSearchIndexTaskResult(progress, stats)
|
|
reportSearchIndexProgress(progress, 100, "Meilisearch is not configured")
|
|
return stats, nil
|
|
}
|
|
stats.Configured = true
|
|
|
|
lock, locked, err := i.events.TryAdvisoryLock(ctx, searchIndexMaintenanceLockID)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
if !locked {
|
|
stats.Skipped = true
|
|
stats.Reason = "another search index maintenance task is running"
|
|
setSearchIndexTaskResult(progress, stats)
|
|
reportSearchIndexProgress(progress, 100, "Another catalog search index maintenance task is already running")
|
|
return stats, nil
|
|
}
|
|
defer lock.Close(context.Background())
|
|
|
|
state, err := i.events.GetState(ctx, SearchProviderMeilisearch)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
if state.ActiveIndexUID == "" || state.SchemaVersion != catalogSearchMeilisearchSchemaVersion(settings.Embedder, settings.IndexTypes, settings.SemanticEnabled, settings.BinaryQuantized) {
|
|
stats.Skipped = true
|
|
stats.Reason = "active search index is missing or stale; run rebuild_catalog_search_index"
|
|
setSearchIndexTaskResult(progress, stats)
|
|
reportSearchIndexProgress(progress, 100, "Catalog search index needs a rebuild")
|
|
return stats, nil
|
|
}
|
|
stats.ActiveIndexUID = state.ActiveIndexUID
|
|
|
|
reportSearchIndexProgress(progress, 0, "Loading pending catalog search events")
|
|
events, err := i.events.ListPending(ctx, SearchProviderMeilisearch, settings.SyncBatchSize)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
if len(events) == 0 {
|
|
stats.DocumentCount = state.DocumentCount
|
|
if vectorCount, err := countCatalogSearchVectorDocuments(ctx, i.pool, settings.IndexTypes, ""); err == nil {
|
|
stats.VectorDocCount = vectorCount
|
|
}
|
|
setSearchIndexTaskResult(progress, stats)
|
|
reportSearchIndexProgress(progress, 100, "Catalog search index is already current")
|
|
return stats, nil
|
|
}
|
|
stats.Events = len(events)
|
|
ids := make([]int64, 0, len(events))
|
|
maxID := int64(0)
|
|
for _, event := range events {
|
|
ids = append(ids, event.ID)
|
|
if event.ID > maxID {
|
|
maxID = event.ID
|
|
}
|
|
}
|
|
|
|
upsertIDs, deleteIDs := coalesceSearchIndexEvents(events)
|
|
reportSearchIndexProgress(progress, 20, "Building changed catalog search documents")
|
|
docs, err := i.LoadDocumentsByIDs(ctx, upsertIDs, settings.IndexTypes, settings.Embedder, settings.SemanticEnabled, settings.BinaryQuantized)
|
|
if err != nil {
|
|
_ = i.events.MarkFailed(ctx, ids, err)
|
|
return stats, err
|
|
}
|
|
found := make(map[string]struct{}, len(docs))
|
|
for _, doc := range docs {
|
|
found[doc.ContentID] = struct{}{}
|
|
}
|
|
for _, id := range upsertIDs {
|
|
if _, ok := found[id]; !ok {
|
|
deleteIDs = append(deleteIDs, id)
|
|
}
|
|
}
|
|
deleteIDs = compactNonEmptyStrings(deleteIDs)
|
|
|
|
if len(deleteIDs) > 0 {
|
|
reportSearchIndexProgress(progress, 45, "Deleting stale catalog search documents")
|
|
taskID, err := client.DeleteDocuments(ctx, state.ActiveIndexUID, deleteIDs)
|
|
if err != nil {
|
|
_ = i.events.MarkFailed(ctx, ids, err)
|
|
return stats, err
|
|
}
|
|
if err := client.WaitTask(ctx, taskID); err != nil {
|
|
_ = i.events.MarkFailed(ctx, ids, err)
|
|
return stats, err
|
|
}
|
|
stats.Deleted = len(deleteIDs)
|
|
}
|
|
if len(docs) > 0 {
|
|
reportSearchIndexProgress(progress, 70, "Upserting catalog search documents")
|
|
for _, batch := range catalogSearchDocumentPayloadBatches(docs, meilisearchMaxDocumentPayloadBytes) {
|
|
taskID, err := client.AddDocuments(ctx, state.ActiveIndexUID, batch)
|
|
if err != nil {
|
|
_ = i.events.MarkFailed(ctx, ids, err)
|
|
return stats, err
|
|
}
|
|
if err := client.WaitTask(ctx, taskID); err != nil {
|
|
_ = i.events.MarkFailed(ctx, ids, err)
|
|
return stats, err
|
|
}
|
|
stats.Upserted += len(batch)
|
|
}
|
|
}
|
|
|
|
docCount, err := client.Stats(ctx, state.ActiveIndexUID)
|
|
if err != nil {
|
|
_ = i.events.MarkFailed(ctx, ids, err)
|
|
return stats, err
|
|
}
|
|
stats.DocumentCount = docCount
|
|
if vectorCount, err := countCatalogSearchVectorDocuments(ctx, i.pool, settings.IndexTypes, ""); err == nil {
|
|
stats.VectorDocCount = vectorCount
|
|
}
|
|
stats.LastProcessedID = maxID
|
|
if err := i.events.MarkProcessed(ctx, ids); err != nil {
|
|
return stats, err
|
|
}
|
|
if err := i.events.UpdateStateAfterSync(ctx, SearchProviderMeilisearch, maxID, docCount); err != nil {
|
|
return stats, err
|
|
}
|
|
setSearchIndexTaskResult(progress, stats)
|
|
reportSearchIndexProgress(progress, 100, fmt.Sprintf("Synced %d catalog search events", len(events)))
|
|
return stats, nil
|
|
}
|
|
|
|
func (i *CatalogSearchIndexer) Rebuild(ctx context.Context, progress SearchIndexProgressReporter) (CatalogSearchIndexRebuildStats, error) {
|
|
stats := CatalogSearchIndexRebuildStats{}
|
|
settings, client, ok, err := i.loadClient(ctx)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
if !ok {
|
|
stats.Skipped = true
|
|
stats.Reason = "meilisearch is not configured"
|
|
setSearchIndexTaskResult(progress, stats)
|
|
reportSearchIndexProgress(progress, 100, "Meilisearch is not configured")
|
|
return stats, nil
|
|
}
|
|
stats.Configured = true
|
|
|
|
lock, locked, err := i.events.TryAdvisoryLock(ctx, searchIndexMaintenanceLockID)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
if !locked {
|
|
stats.Skipped = true
|
|
stats.Reason = "another search index maintenance task is running"
|
|
setSearchIndexTaskResult(progress, stats)
|
|
reportSearchIndexProgress(progress, 100, "Another catalog search index maintenance task is already running")
|
|
return stats, nil
|
|
}
|
|
defer lock.Close(context.Background())
|
|
|
|
rebuildEventHighWater, err := i.events.MaxEventID(ctx, SearchProviderMeilisearch)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
priorState, err := i.events.GetState(ctx, SearchProviderMeilisearch)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
totalDocs, err := countCatalogSearchEligibleDocuments(ctx, i.pool, settings.IndexTypes)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
|
|
buildIndexUID := fmt.Sprintf("%s_rebuild_%d", settings.MeilisearchIndex, time.Now().Unix())
|
|
stats.ActiveIndexUID = buildIndexUID
|
|
reportSearchIndexProgress(progress, 0, "Creating catalog search index")
|
|
taskID, err := client.CreateIndex(ctx, buildIndexUID)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
if err := client.WaitTask(ctx, taskID); err != nil {
|
|
return stats, err
|
|
}
|
|
taskID, err = client.UpdateSettings(ctx, buildIndexUID, catalogSearchMeilisearchSettings(settings.Embedder, settings.SemanticEnabled, settings.BinaryQuantized))
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
if err := client.WaitTask(ctx, taskID); err != nil {
|
|
return stats, err
|
|
}
|
|
|
|
lastID := ""
|
|
queuedTasks := make([]queuedMeilisearchTask, 0, settings.RebuildQueueDepth)
|
|
for {
|
|
docs, err := i.LoadDocumentsAfter(ctx, lastID, settings.RebuildBatchSize, settings.IndexTypes, settings.Embedder, settings.SemanticEnabled, settings.BinaryQuantized)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
if len(docs) == 0 {
|
|
break
|
|
}
|
|
for _, batch := range catalogSearchDocumentPayloadBatches(docs, meilisearchMaxDocumentPayloadBytes) {
|
|
taskID, err := client.AddDocuments(ctx, buildIndexUID, batch)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
queuedTasks = append(queuedTasks, queuedMeilisearchTask{
|
|
task: taskID,
|
|
docCount: len(batch),
|
|
vecCount: catalogSearchVectorDocumentCount(batch),
|
|
})
|
|
submitted := stats.DocumentCount + queuedDocumentCount(queuedTasks)
|
|
reportSearchIndexProgress(progress, rebuildIndexingPercent(submitted, totalDocs), fmt.Sprintf("Submitted %d of %d catalog items", submitted, totalDocs))
|
|
if len(queuedTasks) >= settings.RebuildQueueDepth {
|
|
if err := waitNextMeilisearchTask(ctx, client, &queuedTasks, &stats, progress, totalDocs); err != nil {
|
|
return stats, err
|
|
}
|
|
}
|
|
}
|
|
lastID = docs[len(docs)-1].ContentID
|
|
}
|
|
for len(queuedTasks) > 0 {
|
|
if err := waitNextMeilisearchTask(ctx, client, &queuedTasks, &stats, progress, totalDocs); err != nil {
|
|
return stats, err
|
|
}
|
|
}
|
|
docCount, err := client.Stats(ctx, buildIndexUID)
|
|
if err != nil {
|
|
return stats, err
|
|
}
|
|
stats.DocumentCount = docCount
|
|
if vectorCount, err := countCatalogSearchVectorDocuments(ctx, i.pool, settings.IndexTypes, ""); err == nil {
|
|
stats.VectorDocCount = vectorCount
|
|
}
|
|
// Swap the state pointer BEFORE marking events processed. If the process
|
|
// dies between the two, still-pending events simply replay into the new
|
|
// active index as idempotent upserts on the next sync. The reverse order
|
|
// would mark events processed while the old index is still active, losing
|
|
// those changes from the served index until the next rebuild.
|
|
if err := i.events.UpdateStateAfterRebuild(ctx, SearchProviderMeilisearch, buildIndexUID, catalogSearchMeilisearchSchemaVersion(settings.Embedder, settings.IndexTypes, settings.SemanticEnabled, settings.BinaryQuantized), docCount, rebuildEventHighWater); err != nil {
|
|
return stats, err
|
|
}
|
|
if err := i.events.MarkProcessedThrough(ctx, SearchProviderMeilisearch, rebuildEventHighWater); err != nil {
|
|
return stats, err
|
|
}
|
|
reportSearchIndexProgress(progress, 95, "Removing superseded catalog search indexes")
|
|
removed, err := cleanupSupersededMeilisearchIndexes(ctx, client, settings.MeilisearchIndex, buildIndexUID, priorState.ActiveIndexUID)
|
|
stats.RemovedIndexes = removed
|
|
if err != nil {
|
|
// The new index is already active; a failed cleanup costs disk on the
|
|
// Meilisearch instance, not correctness, and the next rebuild retries.
|
|
slog.WarnContext(ctx, "catalog search: failed to remove superseded meilisearch indexes", "component", "catalog", "err", err, "removed", removed)
|
|
}
|
|
setSearchIndexTaskResult(progress, stats)
|
|
reportSearchIndexProgress(progress, 100, fmt.Sprintf("Rebuilt catalog search index with %d documents", docCount))
|
|
return stats, nil
|
|
}
|
|
|
|
// rebuildIndexingPercent maps rebuild document progress onto the 5-90% band of
|
|
// the task's progress bar (index creation sits below, finalization above).
|
|
func rebuildIndexingPercent(done, total int) float64 {
|
|
if total <= 0 {
|
|
return 50
|
|
}
|
|
if done > total {
|
|
done = total
|
|
}
|
|
return 5 + 85*float64(done)/float64(total)
|
|
}
|
|
|
|
// cleanupSupersededMeilisearchIndexes deletes indexes this rebuild has made
|
|
// unreachable: the previously active index and any `<prefix>_rebuild_*`
|
|
// leftovers from failed or superseded runs. Without it, every rebuild leaks a
|
|
// full copy of the catalog on the Meilisearch instance.
|
|
func cleanupSupersededMeilisearchIndexes(ctx context.Context, client *meilisearchClient, indexPrefix, activeUID, previousActiveUID string) (int, error) {
|
|
uids, err := client.ListIndexUIDs(ctx)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
removed := 0
|
|
for _, uid := range staleCatalogSearchIndexUIDs(uids, indexPrefix, activeUID, previousActiveUID) {
|
|
task, err := client.DeleteIndex(ctx, uid)
|
|
if err != nil {
|
|
return removed, err
|
|
}
|
|
if err := client.WaitTask(ctx, task); err != nil {
|
|
return removed, err
|
|
}
|
|
removed++
|
|
}
|
|
return removed, nil
|
|
}
|
|
|
|
// staleCatalogSearchIndexUIDs selects which index uids a finished rebuild
|
|
// should delete: every `<prefix>_rebuild_` index except the newly active one,
|
|
// plus the previously active index (which may predate the rebuild naming
|
|
// scheme). Indexes outside the prefix are never touched, so a shared
|
|
// Meilisearch instance stays safe.
|
|
func staleCatalogSearchIndexUIDs(uids []string, indexPrefix, activeUID, previousActiveUID string) []string {
|
|
rebuildPrefix := indexPrefix + "_rebuild_"
|
|
var stale []string
|
|
for _, uid := range uids {
|
|
if uid == "" || uid == activeUID {
|
|
continue
|
|
}
|
|
if strings.HasPrefix(uid, rebuildPrefix) || uid == previousActiveUID {
|
|
stale = append(stale, uid)
|
|
}
|
|
}
|
|
return stale
|
|
}
|
|
|
|
func waitNextMeilisearchTask(
|
|
ctx context.Context,
|
|
client *meilisearchClient,
|
|
queue *[]queuedMeilisearchTask,
|
|
stats *CatalogSearchIndexRebuildStats,
|
|
progress SearchIndexProgressReporter,
|
|
totalDocs int,
|
|
) error {
|
|
if len(*queue) == 0 {
|
|
return nil
|
|
}
|
|
next := (*queue)[0]
|
|
*queue = (*queue)[1:]
|
|
if err := client.WaitTask(ctx, next.task); err != nil {
|
|
return err
|
|
}
|
|
stats.DocumentCount += next.docCount
|
|
stats.VectorDocCount += next.vecCount
|
|
reportSearchIndexProgress(progress, rebuildIndexingPercent(stats.DocumentCount, totalDocs), fmt.Sprintf("Indexed %d of %d catalog items", stats.DocumentCount, totalDocs))
|
|
return nil
|
|
}
|
|
|
|
func queuedDocumentCount(queue []queuedMeilisearchTask) int {
|
|
total := 0
|
|
for _, task := range queue {
|
|
total += task.docCount
|
|
}
|
|
return total
|
|
}
|
|
|
|
func catalogSearchDocumentPayloadBatches(docs []catalogSearchDocument, maxBytes int) [][]catalogSearchDocument {
|
|
if len(docs) == 0 {
|
|
return nil
|
|
}
|
|
if maxBytes <= 0 {
|
|
return [][]catalogSearchDocument{docs}
|
|
}
|
|
batches := make([][]catalogSearchDocument, 0, int(math.Ceil(float64(len(docs))/1000)))
|
|
current := make([]catalogSearchDocument, 0, len(docs))
|
|
currentBytes := 2
|
|
for _, doc := range docs {
|
|
docBytes := estimateCatalogSearchDocumentJSONBytes(doc)
|
|
if len(current) > 0 && currentBytes+docBytes+1 > maxBytes {
|
|
batches = append(batches, current)
|
|
current = nil
|
|
currentBytes = 2
|
|
}
|
|
current = append(current, doc)
|
|
currentBytes += docBytes + 1
|
|
}
|
|
if len(current) > 0 {
|
|
batches = append(batches, current)
|
|
}
|
|
return batches
|
|
}
|
|
|
|
func estimateCatalogSearchDocumentJSONBytes(doc catalogSearchDocument) int {
|
|
size := 512 +
|
|
len(doc.ContentID) +
|
|
len(doc.Type) +
|
|
len(doc.Title) +
|
|
len(doc.SortTitle) +
|
|
len(doc.OriginalTitle) +
|
|
len(doc.Overview) +
|
|
len(doc.Tagline)
|
|
size += estimateStringSliceJSONBytes(doc.TitleVariants)
|
|
size += estimateStringSliceJSONBytes(doc.Genres)
|
|
size += estimateStringSliceJSONBytes(doc.Studios)
|
|
size += estimateStringSliceJSONBytes(doc.Networks)
|
|
size += estimateStringSliceJSONBytes(doc.Countries)
|
|
size += estimateStringSliceJSONBytes(doc.Keywords)
|
|
size += estimateStringSliceJSONBytes(doc.People)
|
|
for embedder, vector := range doc.Vectors {
|
|
size += len(embedder) + 64 + len(vector)*24
|
|
}
|
|
return size
|
|
}
|
|
|
|
func estimateStringSliceJSONBytes(values []string) int {
|
|
size := 2
|
|
for _, value := range values {
|
|
size += len(value) + 4
|
|
}
|
|
return size
|
|
}
|
|
|
|
func (i *CatalogSearchIndexer) loadClient(ctx context.Context) (CatalogSearchSettings, *meilisearchClient, bool, error) {
|
|
settings, ok, err := i.loadMeilisearchRuntime(ctx)
|
|
if err != nil || !ok {
|
|
return settings, nil, ok, err
|
|
}
|
|
client, err := newMeilisearchClient(settings.MeilisearchURL, settings.MeilisearchAPIKey, meilisearchIndexingTimeout)
|
|
if err != nil {
|
|
return settings, nil, false, err
|
|
}
|
|
return settings, client, true, nil
|
|
}
|
|
|
|
func (i *CatalogSearchIndexer) loadMeilisearchRuntime(ctx context.Context) (CatalogSearchSettings, bool, error) {
|
|
settings, err := LoadCatalogSearchSettings(ctx, i.settingsStore)
|
|
if err != nil {
|
|
return settings, false, err
|
|
}
|
|
if settings.Provider != SearchProviderMeilisearch || strings.TrimSpace(settings.MeilisearchURL) == "" {
|
|
return settings, false, nil
|
|
}
|
|
return settings, true, nil
|
|
}
|
|
|
|
func (i *CatalogSearchIndexer) CheckConnection(ctx context.Context, settings CatalogSearchSettings) error {
|
|
if settings.Provider != SearchProviderMeilisearch {
|
|
return nil
|
|
}
|
|
client, err := newMeilisearchClient(settings.MeilisearchURL, settings.MeilisearchAPIKey, settings.Timeout)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return client.Health(ctx)
|
|
}
|
|
|
|
func catalogSearchMeilisearchSettings(embedder string, semanticEnabled, binaryQuantized bool) map[string]any {
|
|
settings := map[string]any{
|
|
"displayedAttributes": []string{"content_id", "type"},
|
|
"filterableAttributes": []string{"type", "library_ids"},
|
|
"searchableAttributes": []string{
|
|
"title",
|
|
"original_title",
|
|
"sort_title",
|
|
"title_variants",
|
|
"people",
|
|
"studios",
|
|
"networks",
|
|
"genres",
|
|
"keywords",
|
|
"overview",
|
|
"tagline",
|
|
},
|
|
"pagination": map[string]any{
|
|
"maxTotalHits": meilisearchCandidateScanCap,
|
|
},
|
|
}
|
|
if semanticEnabled {
|
|
embedder, err := NormalizeCatalogSearchEmbedderName(embedder)
|
|
if err != nil {
|
|
embedder = DefaultMeilisearchEmbedder
|
|
}
|
|
settings["embedders"] = catalogSearchMeilisearchEmbedderSettings(embedder, binaryQuantized)
|
|
}
|
|
return settings
|
|
}
|
|
|
|
func coalesceSearchIndexEvents(events []SearchIndexEvent) (upsertIDs []string, deleteIDs []string) {
|
|
ops := make(map[string]string, len(events))
|
|
for _, event := range events {
|
|
switch event.Action {
|
|
case SearchIndexEventRename:
|
|
if event.PreviousContentID != "" {
|
|
ops[event.PreviousContentID] = SearchIndexEventDelete
|
|
}
|
|
if event.ContentID != "" {
|
|
ops[event.ContentID] = SearchIndexEventUpsert
|
|
}
|
|
case SearchIndexEventDelete:
|
|
ops[event.ContentID] = SearchIndexEventDelete
|
|
case SearchIndexEventUpsert:
|
|
ops[event.ContentID] = SearchIndexEventUpsert
|
|
}
|
|
}
|
|
for id, op := range ops {
|
|
switch op {
|
|
case SearchIndexEventDelete:
|
|
deleteIDs = append(deleteIDs, id)
|
|
case SearchIndexEventUpsert:
|
|
upsertIDs = append(upsertIDs, id)
|
|
}
|
|
}
|
|
return compactNonEmptyStrings(upsertIDs), compactNonEmptyStrings(deleteIDs)
|
|
}
|
|
|
|
func reportSearchIndexProgress(progress SearchIndexProgressReporter, percent float64, message string) {
|
|
if progress != nil {
|
|
progress.Report(percent, message)
|
|
}
|
|
}
|
|
|
|
func setSearchIndexTaskResult(progress SearchIndexProgressReporter, result any) {
|
|
if progress == nil {
|
|
return
|
|
}
|
|
data, err := json.Marshal(result)
|
|
if err == nil {
|
|
progress.SetResultData(data)
|
|
}
|
|
}
|
|
|
|
type catalogSearchDocument struct {
|
|
ContentID string `json:"content_id"`
|
|
Type string `json:"type"`
|
|
Title string `json:"title"`
|
|
SortTitle string `json:"sort_title,omitempty"`
|
|
OriginalTitle string `json:"original_title,omitempty"`
|
|
Aliases []string `json:"-"`
|
|
TitleVariants []string `json:"title_variants,omitempty"`
|
|
Year int `json:"year,omitempty"`
|
|
Overview string `json:"overview,omitempty"`
|
|
Tagline string `json:"tagline,omitempty"`
|
|
Genres []string `json:"genres,omitempty"`
|
|
Studios []string `json:"studios,omitempty"`
|
|
Networks []string `json:"networks,omitempty"`
|
|
Countries []string `json:"countries,omitempty"`
|
|
Keywords []string `json:"keywords,omitempty"`
|
|
People []string `json:"people,omitempty"`
|
|
LibraryIDs []int32 `json:"library_ids,omitempty"`
|
|
SchemaVersion int `json:"schema_version"`
|
|
Vectors map[string][]float32 `json:"_vectors,omitempty"`
|
|
}
|
|
|
|
func (i *CatalogSearchIndexer) LoadDocumentsAfter(ctx context.Context, afterContentID string, limit int, itemTypes []string, embedder string, semanticEnabled, binaryQuantized bool) ([]catalogSearchDocument, error) {
|
|
if i == nil || i.pool == nil || limit <= 0 {
|
|
return nil, nil
|
|
}
|
|
typeFilter := normalizeCatalogSearchItemTypes(itemTypes)
|
|
args := []any{afterContentID, limit}
|
|
typeArg := 0
|
|
if len(typeFilter) > 0 {
|
|
typeArg = 3
|
|
args = append(args, typeFilter)
|
|
}
|
|
rows, err := i.pool.Query(ctx, mixedCatalogSearchDocumentSQL(
|
|
`mi.content_id > $1`, `e.content_id > $1`, typeArg, "LIMIT $2"), args...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
docs, err := scanCatalogSearchDocuments(rows)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
setCatalogSearchDocumentSchemaVersion(docs, catalogSearchMeilisearchSchemaVersion(embedder, itemTypes, semanticEnabled, binaryQuantized))
|
|
if err := i.attachDocumentVectors(ctx, docs, embedder, semanticEnabled); err != nil {
|
|
return nil, err
|
|
}
|
|
return docs, nil
|
|
}
|
|
|
|
func (i *CatalogSearchIndexer) LoadDocumentsByIDs(ctx context.Context, contentIDs []string, itemTypes []string, embedder string, semanticEnabled, binaryQuantized bool) ([]catalogSearchDocument, error) {
|
|
contentIDs = compactNonEmptyStrings(contentIDs)
|
|
if i == nil || i.pool == nil || len(contentIDs) == 0 {
|
|
return nil, nil
|
|
}
|
|
typeFilter := normalizeCatalogSearchItemTypes(itemTypes)
|
|
args := []any{contentIDs}
|
|
typeArg := 0
|
|
if len(typeFilter) > 0 {
|
|
typeArg = 2
|
|
args = append(args, typeFilter)
|
|
}
|
|
rows, err := i.pool.Query(ctx, mixedCatalogSearchDocumentSQL(
|
|
`mi.content_id = ANY($1)`, `e.content_id = ANY($1)`, typeArg, ""), args...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
docs, err := scanCatalogSearchDocuments(rows)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
setCatalogSearchDocumentSchemaVersion(docs, catalogSearchMeilisearchSchemaVersion(embedder, itemTypes, semanticEnabled, binaryQuantized))
|
|
if err := i.attachDocumentVectors(ctx, docs, embedder, semanticEnabled); err != nil {
|
|
return nil, err
|
|
}
|
|
return docs, nil
|
|
}
|
|
|
|
// mixedCatalogSearchDocumentSQL first pages narrow IDs across both physical
|
|
// sources, then aggregates people and library memberships only for that batch.
|
|
// This keeps rebuild work proportional to RebuildBatchSize rather than the
|
|
// total number of catalog rows.
|
|
func mixedCatalogSearchDocumentSQL(mediaPredicate, episodePredicate string, typeArg int, limitClause string) string {
|
|
mediaTypePredicate := ""
|
|
episodeTypePredicate := ""
|
|
if typeArg > 0 {
|
|
mediaTypePredicate = fmt.Sprintf(" AND mi.type = ANY($%d)", typeArg)
|
|
episodeTypePredicate = fmt.Sprintf(" AND 'episode' = ANY($%d)", typeArg)
|
|
}
|
|
return fmt.Sprintf(`
|
|
WITH candidates AS (
|
|
SELECT mi.content_id, mi.type
|
|
FROM media_items mi
|
|
WHERE %s
|
|
AND %s%s
|
|
UNION ALL
|
|
SELECT e.content_id, 'episode'::text AS type
|
|
FROM episodes e
|
|
JOIN media_items si ON si.content_id = e.series_id AND si.type = 'series'
|
|
WHERE %s%s
|
|
AND EXISTS (SELECT 1 FROM episode_libraries el WHERE el.episode_id = e.content_id)
|
|
ORDER BY content_id ASC
|
|
%s
|
|
)
|
|
SELECT * FROM (
|
|
SELECT
|
|
mi.content_id,
|
|
mi.type,
|
|
COALESCE(mi.title, ''),
|
|
COALESCE(mi.sort_title, ''),
|
|
COALESCE(mi.original_title, ''),
|
|
COALESCE(aliases.titles, ARRAY[]::text[]),
|
|
COALESCE(mi.year, 0),
|
|
COALESCE(mi.overview, ''),
|
|
COALESCE(mi.tagline, ''),
|
|
COALESCE(mi.genres, ARRAY[]::text[]),
|
|
COALESCE(mi.studios, ARRAY[]::text[]),
|
|
COALESCE(mi.networks, ARRAY[]::text[]),
|
|
COALESCE(mi.countries, ARRAY[]::text[]),
|
|
COALESCE(mi.keywords, ARRAY[]::text[]),
|
|
COALESCE(people.names, ARRAY[]::text[]),
|
|
COALESCE(libraries.ids, ARRAY[]::integer[])
|
|
FROM candidates c
|
|
JOIN media_items mi ON c.type <> 'episode' AND mi.content_id = c.content_id
|
|
LEFT JOIN LATERAL (
|
|
SELECT array_agg(DISTINCT mia.title ORDER BY mia.title) FILTER (WHERE btrim(mia.title) <> '') AS titles
|
|
FROM media_item_aliases mia
|
|
WHERE mia.content_id = mi.content_id
|
|
) aliases ON true
|
|
LEFT JOIN LATERAL (
|
|
SELECT array_agg(DISTINCT p.name) FILTER (WHERE p.name IS NOT NULL AND p.name <> '') AS names
|
|
FROM item_people ip
|
|
JOIN people p ON p.id = ip.person_id
|
|
WHERE ip.content_id = mi.content_id
|
|
) people ON true
|
|
LEFT JOIN LATERAL (
|
|
SELECT array_agg(DISTINCT mil.media_folder_id ORDER BY mil.media_folder_id) AS ids
|
|
FROM media_item_libraries mil
|
|
WHERE mil.content_id = mi.content_id
|
|
) libraries ON true
|
|
UNION ALL
|
|
SELECT
|
|
e.content_id,
|
|
'episode'::text,
|
|
COALESCE(NULLIF(BTRIM(e.title), ''), 'Episode ' || e.episode_number::text),
|
|
COALESCE(NULLIF(BTRIM(e.title), ''), 'Episode ' || e.episode_number::text),
|
|
''::text,
|
|
ARRAY[]::text[],
|
|
COALESCE(EXTRACT(YEAR FROM e.air_date)::integer, 0),
|
|
COALESCE(e.overview, ''),
|
|
''::text,
|
|
ARRAY[]::text[], ARRAY[]::text[], ARRAY[]::text[], ARRAY[]::text[], ARRAY[]::text[], ARRAY[]::text[],
|
|
COALESCE(libraries.ids, ARRAY[]::integer[])
|
|
FROM candidates c
|
|
JOIN episodes e ON c.type = 'episode' AND e.content_id = c.content_id
|
|
LEFT JOIN LATERAL (
|
|
SELECT array_agg(DISTINCT el.media_folder_id ORDER BY el.media_folder_id) AS ids
|
|
FROM episode_libraries el
|
|
WHERE el.episode_id = e.content_id
|
|
) libraries ON true
|
|
) documents
|
|
ORDER BY content_id ASC`,
|
|
mediaPredicate, catalogSearchExcludeMangaChaptersSQL, mediaTypePredicate,
|
|
episodePredicate, episodeTypePredicate, limitClause)
|
|
}
|
|
|
|
func scanCatalogSearchDocuments(rows pgx.Rows) ([]catalogSearchDocument, error) {
|
|
var docs []catalogSearchDocument
|
|
for rows.Next() {
|
|
var doc catalogSearchDocument
|
|
if err := rows.Scan(
|
|
&doc.ContentID,
|
|
&doc.Type,
|
|
&doc.Title,
|
|
&doc.SortTitle,
|
|
&doc.OriginalTitle,
|
|
&doc.Aliases,
|
|
&doc.Year,
|
|
&doc.Overview,
|
|
&doc.Tagline,
|
|
&doc.Genres,
|
|
&doc.Studios,
|
|
&doc.Networks,
|
|
&doc.Countries,
|
|
&doc.Keywords,
|
|
&doc.People,
|
|
&doc.LibraryIDs,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
doc.SchemaVersion = SearchMeilisearchSchemaVersion
|
|
doc.TitleVariants = catalogSearchTitleVariants(doc)
|
|
docs = append(docs, doc)
|
|
}
|
|
return docs, rows.Err()
|
|
}
|
|
|
|
func (i *CatalogSearchIndexer) attachDocumentVectors(ctx context.Context, docs []catalogSearchDocument, embedder string, semanticEnabled bool) error {
|
|
if !semanticEnabled || i == nil || i.pool == nil || len(docs) == 0 {
|
|
return nil
|
|
}
|
|
embedder, err := NormalizeCatalogSearchEmbedderName(embedder)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
ids := make([]string, 0, len(docs))
|
|
for _, doc := range docs {
|
|
if doc.Type != "episode" && strings.TrimSpace(doc.ContentID) != "" {
|
|
ids = append(ids, doc.ContentID)
|
|
}
|
|
}
|
|
var vectors map[string][]float32
|
|
if len(ids) > 0 {
|
|
if vectors, err = loadCatalogSearchVectors(ctx, i.pool, ids); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
// Always run the per-document pass: even an episode-only batch needs the
|
|
// explicit `_vectors.<embedder>: null` opt-out on every document.
|
|
setCatalogSearchDocumentVectors(docs, vectors, embedder)
|
|
return nil
|
|
}
|
|
|
|
func loadCatalogSearchVectors(ctx context.Context, pool *pgxpool.Pool, contentIDs []string) (map[string][]float32, error) {
|
|
contentIDs = compactNonEmptyStrings(contentIDs)
|
|
if pool == nil || len(contentIDs) == 0 {
|
|
return nil, nil
|
|
}
|
|
rows, err := pool.Query(ctx, `
|
|
SELECT media_item_id, embedding
|
|
FROM media_item_embeddings
|
|
WHERE media_item_id = ANY($1)
|
|
`, contentIDs)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("load catalog search vectors: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
|
|
vectors := make(map[string][]float32, len(contentIDs))
|
|
for rows.Next() {
|
|
var contentID string
|
|
var vector pgvector.Vector
|
|
if err := rows.Scan(&contentID, &vector); err != nil {
|
|
return nil, fmt.Errorf("scan catalog search vector: %w", err)
|
|
}
|
|
canonical, err := embeddingvectors.EnsureCanonicalDimensions(vector.Slice())
|
|
if err != nil {
|
|
return nil, fmt.Errorf("canonicalize catalog search vector for %s: %w", contentID, err)
|
|
}
|
|
vectors[contentID] = canonical
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, fmt.Errorf("iterate catalog search vectors: %w", err)
|
|
}
|
|
return vectors, nil
|
|
}
|
|
|
|
func setCatalogSearchDocumentVectors(docs []catalogSearchDocument, vectors map[string][]float32, embedder string) int {
|
|
if len(docs) == 0 || strings.TrimSpace(embedder) == "" {
|
|
return 0
|
|
}
|
|
count := 0
|
|
for idx := range docs {
|
|
if docs[idx].Type == "episode" {
|
|
// Episodes are keyword-only, but a userProvided embedder requires
|
|
// every document to either supply vectors or opt out explicitly
|
|
// with `_vectors.<embedder>: null`; omitting _vectors entirely
|
|
// fails the whole indexing task.
|
|
docs[idx].Vectors = map[string][]float32{embedder: nil}
|
|
continue
|
|
}
|
|
vector := vectors[docs[idx].ContentID]
|
|
docs[idx].Vectors = map[string][]float32{embedder: nil}
|
|
if len(vector) > 0 {
|
|
docs[idx].Vectors[embedder] = vector
|
|
count++
|
|
}
|
|
}
|
|
return count
|
|
}
|
|
|
|
func setCatalogSearchDocumentSchemaVersion(docs []catalogSearchDocument, schemaVersion int) {
|
|
for idx := range docs {
|
|
docs[idx].SchemaVersion = schemaVersion
|
|
}
|
|
}
|
|
|
|
func catalogSearchVectorDocumentCount(docs []catalogSearchDocument) int {
|
|
count := 0
|
|
for _, doc := range docs {
|
|
for _, vector := range doc.Vectors {
|
|
if len(vector) > 0 {
|
|
count++
|
|
break
|
|
}
|
|
}
|
|
}
|
|
return count
|
|
}
|
|
|
|
// countCatalogSearchVectorDocuments returns the total number of embed-eligible
|
|
// items carrying a current-model embedding across the requested item types
|
|
// (nil/empty => all types). model="" counts every model. This is the numerator
|
|
// of catalogSemanticCoverageByType summed across types; it applies the same
|
|
// embed-eligibility predicate so it never counts vectors on ineligible
|
|
// (unmatched, non-book) items.
|
|
func countCatalogSearchVectorDocuments(ctx context.Context, q coverageQuerier, itemTypes []string, model string) (int, error) {
|
|
if q == nil {
|
|
return 0, nil
|
|
}
|
|
// Shares the eligibility + model + type predicate with
|
|
// semanticCoverageVectorizedByTypeSQL (the per-type numerator); this is the
|
|
// same count summed across types. $1 is always referenced (the explicit
|
|
// ::text[] cast lets Postgres infer the type when it is NULL), so a nil type
|
|
// filter does not trip an "undetermined parameter" error.
|
|
var typeArg any
|
|
if typeFilter := normalizeCatalogSearchItemTypes(itemTypes); len(typeFilter) > 0 {
|
|
typeArg = typeFilter
|
|
}
|
|
var count int
|
|
if err := q.QueryRow(ctx, `
|
|
SELECT COUNT(*)
|
|
FROM media_item_embeddings e
|
|
JOIN media_items mi ON mi.content_id = e.media_item_id
|
|
WHERE `+catalogSearchExcludeMangaChaptersSQL+`
|
|
AND ($1::text[] IS NULL OR mi.type = ANY($1))
|
|
AND (mi.status = 'matched' OR mi.type IN ('audiobook','ebook'))
|
|
AND ($2 = '' OR e.model = $2)
|
|
`, typeArg, model).Scan(&count); err != nil {
|
|
return 0, fmt.Errorf("count catalog search vector documents: %w", err)
|
|
}
|
|
return count, nil
|
|
}
|
|
|
|
// countCatalogSearchEligibleDocuments counts the items a rebuild will index
|
|
// (nil/empty itemTypes => all types). It applies the same predicate as the
|
|
// document loaders — NOT the vector-eligibility predicate — so the total is an
|
|
// exact denominator for rebuild progress reporting.
|
|
func countCatalogSearchEligibleDocuments(ctx context.Context, q coverageQuerier, itemTypes []string) (int, error) {
|
|
if q == nil {
|
|
return 0, nil
|
|
}
|
|
var typeArg any
|
|
if typeFilter := normalizeCatalogSearchItemTypes(itemTypes); len(typeFilter) > 0 {
|
|
typeArg = typeFilter
|
|
}
|
|
var count int
|
|
if err := q.QueryRow(ctx, `
|
|
SELECT
|
|
(SELECT COUNT(*)
|
|
FROM media_items mi
|
|
WHERE `+catalogSearchExcludeMangaChaptersSQL+`
|
|
AND ($1::text[] IS NULL OR mi.type = ANY($1)))
|
|
+
|
|
(SELECT COUNT(*)
|
|
FROM episodes e
|
|
JOIN media_items si ON si.content_id = e.series_id AND si.type = 'series'
|
|
WHERE ($1::text[] IS NULL OR 'episode' = ANY($1))
|
|
AND EXISTS (SELECT 1 FROM episode_libraries el WHERE el.episode_id = e.content_id))
|
|
`, typeArg).Scan(&count); err != nil {
|
|
return 0, fmt.Errorf("count catalog search eligible documents: %w", err)
|
|
}
|
|
return count, nil
|
|
}
|
|
|
|
func catalogSearchTitleVariants(doc catalogSearchDocument) []string {
|
|
variants := []string{
|
|
doc.Title,
|
|
doc.SortTitle,
|
|
doc.OriginalTitle,
|
|
normalizeTitleForComparison(doc.Title),
|
|
normalizeTitleForComparison(doc.SortTitle),
|
|
normalizeTitleForComparison(doc.OriginalTitle),
|
|
}
|
|
for _, alias := range doc.Aliases {
|
|
variants = append(variants, alias, normalizeTitleForComparison(alias))
|
|
}
|
|
return compactNonEmptyStrings(variants)
|
|
}
|