mirror of
https://github.com/opencloud-eu/opencloud.git
synced 2026-08-01 02:11:15 -04:00
explore(search): bleve reset+rescan instead of in-place migration
Alternative that drops the in-place rebuild machinery (buildReplacement, swap, recoverMigration, verifyMigrated, the fixup hook) and just wipes the outdated bleve index so a full 'search index --force-rescan' repopulates it. ~272 lines lighter, but trades the fast no-rescan migration for a Tika re-crawl on every schema bump. Kept in history for reference; reverted in the next commit - the in-place migration is worth the ~300 lines for bleve too.
This commit is contained in:
@@ -3,7 +3,6 @@ package bleve
|
|||||||
import (
|
import (
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io/fs"
|
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strconv"
|
"strconv"
|
||||||
@@ -15,235 +14,66 @@ import (
|
|||||||
"github.com/opencloud-eu/opencloud/services/search/pkg/search"
|
"github.com/opencloud-eu/opencloud/services/search/pkg/search"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const revisionKey = "_schema_revision"
|
||||||
revisionKey = "_schema_revision"
|
|
||||||
migratingSuffix = ".migrating"
|
|
||||||
backupSuffix = ".bak"
|
|
||||||
migrateBatch = 100
|
|
||||||
migrateLogInterval = 50000 // how often buildReplacement logs progress, in documents
|
|
||||||
)
|
|
||||||
|
|
||||||
// OpenOrMigrate opens the index and, when autoMigrate is set and the stored
|
// OpenOrMigrate opens the index and, when autoMigrate is set and the stored
|
||||||
// schema is incompatible, rebuilds it in place first. The migration triggers on
|
// schema is incompatible, resets it first so a re-crawl can repopulate it at the
|
||||||
// a schema revision bump, not on a mapping diff, so a breaking change without a
|
// current schema. The reset triggers on a schema revision bump, not on a mapping
|
||||||
// matching bump still refuses to start. Returns the number of documents migrated
|
// diff, so a breaking change without a matching bump still refuses to start.
|
||||||
// (0 for none). Opening first means a fresh install and the common up-to-date
|
|
||||||
// case open the index only once; the migration path is taken only on a mismatch.
|
|
||||||
func OpenOrMigrate(root string, autoMigrate bool, logger log.Logger) (bleve.Index, searchmapping.Classification, int, error) {
|
func OpenOrMigrate(root string, autoMigrate bool, logger log.Logger) (bleve.Index, searchmapping.Classification, int, error) {
|
||||||
idx, classification, openErr := NewIndex(root)
|
idx, classification, openErr := NewIndex(root)
|
||||||
if openErr == nil || !autoMigrate || !errors.Is(openErr, searchmapping.ErrManualActionRequired) {
|
if openErr == nil || !autoMigrate || !errors.Is(openErr, searchmapping.ErrManualActionRequired) {
|
||||||
// opened cleanly, or auto-migrate is off, or the failure is not a
|
|
||||||
// schema mismatch we could migrate away
|
|
||||||
return idx, classification, 0, openErr
|
return idx, classification, 0, openErr
|
||||||
}
|
}
|
||||||
|
|
||||||
// the stored schema is incompatible; rebuild only if the revision was bumped,
|
// incompatible schema: MigrateIndex drops a bumped-revision index so the
|
||||||
// otherwise surface the original error (a schema change without a bump)
|
// reopen below rebuilds it empty; a breaking change without a bump is a no-op
|
||||||
migrated, err := MigrateIndex(root, logger)
|
// there and is surfaced by that reopen instead
|
||||||
|
dropped, err := MigrateIndex(root, logger)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, classification, 0, err
|
return nil, classification, 0, err
|
||||||
}
|
}
|
||||||
if migrated == 0 {
|
|
||||||
return nil, classification, 0, openErr
|
|
||||||
}
|
|
||||||
idx, classification, err = NewIndex(root)
|
idx, classification, err = NewIndex(root)
|
||||||
return idx, classification, migrated, err
|
return idx, classification, dropped, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// MigrateIndex rebuilds the index from the documents it already holds when the
|
// MigrateIndex resets the index when its stored schema revision is older than
|
||||||
// stored schema revision is older than search.SchemaRevision, else a no-op. It
|
// search.SchemaRevision: it drops the outdated index so a full re-crawl
|
||||||
// takes the exclusive bolt lock, so the search service must be stopped.
|
// ('opencloud search index --all-spaces --force-rescan') repopulates it at the
|
||||||
|
// current schema. No-op when the index is missing or already current. It takes
|
||||||
|
// the exclusive bolt lock, so the search service must be stopped. Returns the
|
||||||
|
// number of documents that were dropped.
|
||||||
func MigrateIndex(root string, logger log.Logger) (int, error) {
|
func MigrateIndex(root string, logger log.Logger) (int, error) {
|
||||||
dest := filepath.Join(root, "bleve")
|
dest := filepath.Join(root, "bleve")
|
||||||
tmp := dest + migratingSuffix
|
|
||||||
bak := dest + backupSuffix
|
|
||||||
|
|
||||||
if err := recoverMigration(root); err != nil {
|
idx, err := bleve.OpenUsing(dest, openRuntimeConfig)
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
|
|
||||||
old, err := bleve.OpenUsing(dest, openRuntimeConfig)
|
|
||||||
if errors.Is(err, bleve.ErrorIndexPathDoesNotExist) {
|
if errors.Is(err, bleve.ErrorIndexPathDoesNotExist) {
|
||||||
return 0, nil // nothing to migrate yet; NewIndex creates a fresh index
|
return 0, nil // fresh install, nothing to migrate
|
||||||
}
|
}
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return 0, fmt.Errorf("open index (is the search service still running?): %w", err)
|
return 0, fmt.Errorf("open index (is the search service still running?): %w", err)
|
||||||
}
|
}
|
||||||
|
stored, rerr := readRevision(idx)
|
||||||
stored, err := readRevision(old)
|
count, cerr := idx.DocCount()
|
||||||
if err != nil {
|
_ = idx.Close()
|
||||||
_ = old.Close()
|
switch {
|
||||||
return 0, err
|
case rerr != nil:
|
||||||
}
|
return 0, rerr
|
||||||
if stored >= search.SchemaRevision {
|
case cerr != nil:
|
||||||
_ = old.Close()
|
return 0, cerr
|
||||||
return 0, nil
|
case stored >= search.SchemaRevision:
|
||||||
|
return 0, nil // already current
|
||||||
}
|
}
|
||||||
|
|
||||||
total, err := old.DocCount()
|
if err := os.RemoveAll(dest); err != nil {
|
||||||
if err != nil {
|
|
||||||
_ = old.Close()
|
|
||||||
return 0, err
|
return 0, err
|
||||||
}
|
}
|
||||||
logger.Info().Uint64("documents", total).Msg("starting search index migration")
|
logger.Warn().Msgf("the bleve index at %s was reset for schema revision %d (%d documents dropped); run 'opencloud search index --all-spaces --force-rescan' to repopulate", root, search.SchemaRevision, count)
|
||||||
|
return int(count), nil
|
||||||
copied, err := buildReplacement(old, tmp, total, logger)
|
|
||||||
if err != nil {
|
|
||||||
_ = old.Close()
|
|
||||||
_ = os.RemoveAll(tmp)
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
|
|
||||||
// verify before the swap, so a failure leaves the original untouched
|
|
||||||
if err := verifyMigrated(tmp, total); err != nil {
|
|
||||||
_ = old.Close()
|
|
||||||
_ = os.RemoveAll(tmp)
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := swap(old, dest, tmp, bak); err != nil {
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
logger.Info().Int("documents", copied).Msg("search index migration complete")
|
|
||||||
return copied, nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// buildReplacement copies every document from src into a fresh index at tmp via
|
// writeRevision stamps the current revision so a fresh index is not seen as
|
||||||
// the same Deserialize -> PrepareForIndex round-trip the Move/Delete/Restore
|
// outdated on the next start.
|
||||||
// paths use, and stamps the current revision.
|
|
||||||
func buildReplacement(src bleve.Index, tmp string, total uint64, logger log.Logger) (int, error) {
|
|
||||||
if err := os.RemoveAll(tmp); err != nil {
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
m, err := NewMapping()
|
|
||||||
if err != nil {
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
dst, err := bleve.New(tmp, m)
|
|
||||||
if err != nil {
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
|
|
||||||
batch, err := NewBatch(dst, migrateBatch)
|
|
||||||
if err != nil {
|
|
||||||
_ = dst.Close()
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
|
|
||||||
var copied int
|
|
||||||
nextLog := migrateLogInterval
|
|
||||||
var after []string
|
|
||||||
for {
|
|
||||||
req := bleve.NewSearchRequest(bleve.NewMatchAllQuery())
|
|
||||||
req.Size = migrateBatch
|
|
||||||
req.Fields = []string{"*"}
|
|
||||||
req.SortBy([]string{"_id"}) // total order, no skips or dupes across pages
|
|
||||||
req.SearchAfter = after
|
|
||||||
res, err := src.Search(req)
|
|
||||||
if err != nil {
|
|
||||||
_ = dst.Close()
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
if len(res.Hits) == 0 {
|
|
||||||
break
|
|
||||||
}
|
|
||||||
for _, hit := range res.Hits {
|
|
||||||
r := legacyHitToResource(hit.Fields)
|
|
||||||
if err := batch.Upsert(hit.ID, r); err != nil { // hit.ID: the bleve key is authoritative
|
|
||||||
_ = dst.Close()
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
copied++
|
|
||||||
}
|
|
||||||
if copied >= nextLog {
|
|
||||||
logger.Info().Msgf("migrating search index: %d of %d documents", copied, total)
|
|
||||||
nextLog += migrateLogInterval
|
|
||||||
}
|
|
||||||
after = []string{res.Hits[len(res.Hits)-1].ID}
|
|
||||||
}
|
|
||||||
if err := batch.Push(); err != nil {
|
|
||||||
_ = dst.Close()
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
if err := writeRevision(dst); err != nil {
|
|
||||||
_ = dst.Close()
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
return copied, dst.Close()
|
|
||||||
}
|
|
||||||
|
|
||||||
// legacyHitToResource is the fixup hook before re-indexing, symmetric to the
|
|
||||||
// opensearch side. bleve Deserialize is fail-soft, so nothing is stripped here.
|
|
||||||
func legacyHitToResource(fields map[string]any) search.Resource {
|
|
||||||
return *searchmapping.Deserialize[search.Resource](fields)
|
|
||||||
}
|
|
||||||
|
|
||||||
// verifyMigrated checks the rebuilt index holds want documents before the swap
|
|
||||||
// replaces the original. The count catches a pagination drop or batch-push loss;
|
|
||||||
// the mapping needs no check, it is built deterministically from NewMapping().
|
|
||||||
func verifyMigrated(path string, want uint64) error {
|
|
||||||
idx, err := bleve.OpenUsing(path, openRuntimeConfig)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
defer func() { _ = idx.Close() }()
|
|
||||||
|
|
||||||
got, err := idx.DocCount()
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if got != want {
|
|
||||||
return fmt.Errorf("migrated index holds %d documents, expected %d", got, want)
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// swap replaces the live index with the rebuilt one. old still holds the lock
|
|
||||||
// and is closed here right before the rename.
|
|
||||||
func swap(old bleve.Index, dest, tmp, bak string) error {
|
|
||||||
if err := old.Close(); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if err := os.Rename(dest, bak); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if err := os.Rename(tmp, dest); err != nil {
|
|
||||||
_ = os.Rename(bak, dest) // best effort roll back
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
return os.RemoveAll(bak)
|
|
||||||
}
|
|
||||||
|
|
||||||
// recoverMigration cleans up a crashed MigrateIndex and must run before the
|
|
||||||
// index is opened, else NewIndex would recreate an empty index mid-swap. Rolls
|
|
||||||
// back, never forward: a truncated tmp is a valid index, completeness is not
|
|
||||||
// observable.
|
|
||||||
func recoverMigration(root string) error {
|
|
||||||
dest := filepath.Join(root, "bleve")
|
|
||||||
tmp := dest + migratingSuffix
|
|
||||||
bak := dest + backupSuffix
|
|
||||||
|
|
||||||
if _, err := os.Stat(dest); err == nil {
|
|
||||||
return errors.Join(os.RemoveAll(tmp), os.RemoveAll(bak)) // index present: tmp/bak are leftovers
|
|
||||||
} else if !errors.Is(err, fs.ErrNotExist) {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
if _, err := os.Stat(bak); err == nil {
|
|
||||||
// crashed mid-swap: roll back
|
|
||||||
if err := os.RemoveAll(tmp); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
return os.Rename(bak, dest)
|
|
||||||
} else if !errors.Is(err, fs.ErrNotExist) {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
return os.RemoveAll(tmp) // fresh install, or index removed by hand
|
|
||||||
}
|
|
||||||
|
|
||||||
// writeRevision stamps the current revision so a fresh or migrated index is not
|
|
||||||
// seen as outdated.
|
|
||||||
func writeRevision(index bleve.Index) error {
|
func writeRevision(index bleve.Index) error {
|
||||||
return index.SetInternal([]byte(revisionKey), []byte(strconv.Itoa(search.SchemaRevision)))
|
return index.SetInternal([]byte(revisionKey), []byte(strconv.Itoa(search.SchemaRevision)))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,89 +0,0 @@
|
|||||||
package bleve_test
|
|
||||||
|
|
||||||
import (
|
|
||||||
"path/filepath"
|
|
||||||
"sort"
|
|
||||||
"testing"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
bleveSearch "github.com/blevesearch/bleve/v2"
|
|
||||||
libregraph "github.com/opencloud-eu/libre-graph-api-go"
|
|
||||||
"github.com/stretchr/testify/require"
|
|
||||||
|
|
||||||
"github.com/opencloud-eu/opencloud/pkg/conversions"
|
|
||||||
"github.com/opencloud-eu/opencloud/pkg/log"
|
|
||||||
"github.com/opencloud-eu/opencloud/services/search/pkg/bleve"
|
|
||||||
"github.com/opencloud-eu/opencloud/services/search/pkg/content"
|
|
||||||
searchmapping "github.com/opencloud-eu/opencloud/services/search/pkg/mapping"
|
|
||||||
"github.com/opencloud-eu/opencloud/services/search/pkg/search"
|
|
||||||
)
|
|
||||||
|
|
||||||
// fullyPopulated exercises every field, including the facets and slices the
|
|
||||||
// count-only test skipped.
|
|
||||||
func fullyPopulated() search.Resource {
|
|
||||||
return search.Resource{
|
|
||||||
ID: "1$2!3",
|
|
||||||
RootID: "1$2!2",
|
|
||||||
ParentID: "1$2!2",
|
|
||||||
Path: "./a/b/song.mp3",
|
|
||||||
Type: 2,
|
|
||||||
Deleted: true,
|
|
||||||
Hidden: true,
|
|
||||||
Document: content.Document{
|
|
||||||
Title: "the title",
|
|
||||||
Name: "song.mp3",
|
|
||||||
Content: "some extracted content",
|
|
||||||
Size: 123456,
|
|
||||||
Mtime: conversions.ToPointer(time.Date(2026, 3, 15, 12, 0, 0, 0, time.UTC)),
|
|
||||||
MimeType: "audio/mpeg",
|
|
||||||
Tags: []string{"alpha", "beta", "gamma"},
|
|
||||||
Favorites: []string{"user-a", "user-b"},
|
|
||||||
Audio: &libregraph.Audio{
|
|
||||||
Album: libregraph.PtrString("the album"),
|
|
||||||
Artist: libregraph.PtrString("the artist"),
|
|
||||||
Track: libregraph.PtrInt32(7),
|
|
||||||
Year: libregraph.PtrInt32(1998),
|
|
||||||
HasDrm: libregraph.PtrBool(false),
|
|
||||||
},
|
|
||||||
Image: &libregraph.Image{Width: libregraph.PtrInt32(1920), Height: libregraph.PtrInt32(1080)},
|
|
||||||
Photo: &libregraph.Photo{CameraMake: libregraph.PtrString("Canon"), Iso: libregraph.PtrInt32(400), FNumber: libregraph.PtrFloat64(2.8)},
|
|
||||||
Location: &libregraph.GeoCoordinates{
|
|
||||||
Longitude: libregraph.PtrFloat64(11.103870357204285),
|
|
||||||
Latitude: libregraph.PtrFloat64(49.48675890884328),
|
|
||||||
Altitude: libregraph.PtrFloat64(300.0),
|
|
||||||
},
|
|
||||||
},
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestMigratePreservesAllFields(t *testing.T) {
|
|
||||||
root := t.TempDir()
|
|
||||||
want := fullyPopulated()
|
|
||||||
|
|
||||||
old, err := bleveSearch.New(filepath.Join(root, "bleve"), oldMainMapping(t))
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.NoError(t, old.Index(want.ID, want))
|
|
||||||
require.NoError(t, old.Close())
|
|
||||||
|
|
||||||
_, err = bleve.MigrateIndex(root, log.NopLogger())
|
|
||||||
require.NoError(t, err)
|
|
||||||
|
|
||||||
idx, _, err := bleve.NewIndex(root)
|
|
||||||
require.NoError(t, err)
|
|
||||||
defer func() { _ = idx.Close() }()
|
|
||||||
|
|
||||||
// read the migrated document back through the production reconstruction path
|
|
||||||
req := bleveSearch.NewSearchRequest(bleveSearch.NewMatchAllQuery())
|
|
||||||
req.Fields = []string{"*"}
|
|
||||||
res, err := idx.Search(req)
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.Len(t, res.Hits, 1)
|
|
||||||
got := searchmapping.Deserialize[search.Resource](res.Hits[0].Fields)
|
|
||||||
|
|
||||||
require.NotNil(t, got.Mtime, "Mtime")
|
|
||||||
require.True(t, want.Mtime.Equal(*got.Mtime), "Mtime: want %v got %v", want.Mtime, got.Mtime)
|
|
||||||
got.Mtime = want.Mtime // normalize time.Time location/monotonic
|
|
||||||
sort.Strings(got.Tags) // bleve may return multi-value fields reordered
|
|
||||||
sort.Strings(got.Favorites) //
|
|
||||||
require.Equal(t, want, *got)
|
|
||||||
}
|
|
||||||
@@ -7,7 +7,6 @@ import (
|
|||||||
|
|
||||||
bleveSearch "github.com/blevesearch/bleve/v2"
|
bleveSearch "github.com/blevesearch/bleve/v2"
|
||||||
bleveMapping "github.com/blevesearch/bleve/v2/mapping"
|
bleveMapping "github.com/blevesearch/bleve/v2/mapping"
|
||||||
"github.com/blevesearch/bleve/v2/search/query"
|
|
||||||
libregraph "github.com/opencloud-eu/libre-graph-api-go"
|
libregraph "github.com/opencloud-eu/libre-graph-api-go"
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
|
|
||||||
@@ -129,35 +128,23 @@ func TestMigrateIndex(t *testing.T) {
|
|||||||
_, _, err = bleve.NewIndex(root)
|
_, _, err = bleve.NewIndex(root)
|
||||||
require.ErrorIs(t, err, searchmapping.ErrManualActionRequired)
|
require.ErrorIs(t, err, searchmapping.ErrManualActionRequired)
|
||||||
|
|
||||||
// 3) migrate
|
// 3) migrate resets the outdated index, dropping its documents
|
||||||
n, err := bleve.MigrateIndex(root, log.NopLogger())
|
n, err := bleve.MigrateIndex(root, log.NopLogger())
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
require.Equal(t, 3, n)
|
require.Equal(t, 3, n)
|
||||||
|
|
||||||
// 4) the migrated index classifies as equal and holds every document
|
// 4) the reset index classifies as equal and is empty, awaiting a re-crawl
|
||||||
migrated, classification, err := bleve.NewIndex(root)
|
migrated, classification, err := bleve.NewIndex(root)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
require.Equal(t, searchmapping.VerdictEqual, classification.Verdict)
|
require.Equal(t, searchmapping.VerdictEqual, classification.Verdict)
|
||||||
|
|
||||||
count, err := migrated.DocCount()
|
count, err := migrated.DocCount()
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
require.Equal(t, uint64(3), count)
|
require.Equal(t, uint64(0), count, "reset drops the documents; a re-crawl repopulates")
|
||||||
|
|
||||||
// 4a) the trashed document survived (a rescan would have dropped it)
|
|
||||||
trashed, err := migrated.Document("1$1!3")
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.NotNil(t, trashed)
|
|
||||||
|
|
||||||
// 5) location_geopoint was synthesized during migration (absent in the old index)
|
|
||||||
near := query.NewGeoDistanceQuery(11.103870357204285, 49.48675890884328, "1km")
|
|
||||||
near.SetField("location" + searchmapping.GeopointSuffix)
|
|
||||||
res, err := migrated.Search(bleveSearch.NewSearchRequest(near))
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.Equal(t, uint64(3), res.Total, "all three docs match a geo-distance query on the new sibling field")
|
|
||||||
|
|
||||||
require.NoError(t, migrated.Close()) // release the lock before the second run
|
require.NoError(t, migrated.Close()) // release the lock before the second run
|
||||||
|
|
||||||
// 6) idempotent: a second run is a no-op because the revision is now current
|
// 5) idempotent: a second run is a no-op because the revision is now current
|
||||||
n, err = bleve.MigrateIndex(root, log.NopLogger())
|
n, err = bleve.MigrateIndex(root, log.NopLogger())
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
require.Equal(t, 0, n)
|
require.Equal(t, 0, n)
|
||||||
|
|||||||
Reference in New Issue
Block a user