diff --git a/docs/services/search.md b/docs/services/search.md index 6d35dfe..f34767c 100644 --- a/docs/services/search.md +++ b/docs/services/search.md @@ -135,6 +135,15 @@ A paginated endpoint also needs to know how many documents matched. An engine th `beachcomber.Query.QueryByWeights` ranks the fields of `QueryBy`, one weight per field in the same order, as Scout's `query_by_weights` option does. +## Reindexing + +`beachcomber.Service.Sync` pushes one model on demand; a reindex command loops over the table with it. Two helpers cover the index itself: + +- `beachcomber.DropIndex` drops an index and reports whether it existed, so the command can print a different message for an index that was already absent. It uses the optional `beachcomber.IndexDropper` interface; for an engine without it, it calls `Flush` and reports that the index existed. +- `beachcomber.EnsureIndex` creates the index, empty, from the schema a `beachcomber.IndexSchemaProvider` model supplies. Call it before the loop: an `Upsert` of zero documents creates nothing, so without it an empty table leaves no index behind. It uses the optional `beachcomber.IndexEnsurer` interface and does nothing for other engines. + +The Typesense engine implements both. + ## Typesense The Typesense engine follows the Scout Typesense wire contract, so indexes built by a WinterCMS application can be searched by the port: diff --git a/modules/beachcomber/README.md b/modules/beachcomber/README.md index 7530d77..b020846 100644 --- a/modules/beachcomber/README.md +++ b/modules/beachcomber/README.md @@ -20,12 +20,14 @@ The package itself knows no search server. An engine package registers itself fr - The `beachcomber.Searchable` model contract: `SearchableAs` (the index name, prefixed with `search.prefix`), `ToSearchableArray(ctx, db)` (the document, built from the committed row and free to query related rows) and `ShouldBeSearchable`. A model can also implement `beachcomber.IndexSchemaProvider` (the schema the engine creates a missing index with) and `beachcomber.SearchKeyer` (a document key other than the decimal primary key). - The `beachcomber.Engine` driver contract: `Name`, `Configured`, `Upsert`, `Delete`, `Flush` and `SearchIDs` with a `beachcomber.Query`. `beachcomber.Query.QueryByWeights` ranks the `QueryBy` fields, one weight per field. - Search totals: the optional `beachcomber.PageSearcher` interface returns a `beachcomber.SearchResult`, one page of ids plus `Found`, the number of documents the engine matched. Call it through `beachcomber.SearchPage`, which falls back to `SearchIDs` (with `Found` set to the number of ids) for an engine that does not implement it; the `null` engine returns no ids and a zero count. +- Index lifecycle for reindex tooling: the optional `beachcomber.IndexDropper` interface drops an index and reports whether it existed, so a command can tell a dropped index from one that was already absent; call it through `beachcomber.DropIndex`, which falls back to `Flush` (reporting existed) for other engines. The optional `beachcomber.IndexEnsurer` creates an empty index from its schema; `beachcomber.EnsureIndex` calls it and is a no-op for other engines. A full reindex calls it first because an `Upsert` of zero documents creates nothing. - GORM callbacks `beachcomber.CallbackAfterCreate`, `beachcomber.CallbackAfterUpdate` and `beachcomber.CallbackAfterDelete`, installed once per `*gorm.DB`, register the sync with `lagoon.AfterCommit`. - `beachcomber.Service.Sync` and `beachcomber.Service.Remove` run the same gated path on demand, for reindex tooling, and return the error instead of logging it. - Typesense engine (`typesense.Engine`, engine name `typesense`): - Every request carries the `X-TYPESENSE-API-KEY` header. - `Upsert` reads the collection and creates it from the schema on 404. A 409 on create counts as success, and a model without a schema gets an auto-typed collection. It then imports the documents as JSON lines (`Content-Type: text/plain`) with `action=upsert`. Typesense answers 200 even when a document fails, so every answer line is checked and any `"success":false` line is an error. - `Delete` and `Flush` treat 404 as success. + - `DropIndex` sends `DELETE /collections/{index}` and reports existed true on 2xx and false on 404; `EnsureIndex` runs the same read-then-create step as `Upsert`, so a second call sends no create. - `SearchIDs` and `SearchPage` send one request with `q` (default `*`), `query_by`, `query_by_weights`, `filter_by`, `sort_by`, `page` and `per_page`, and return `hits[].document.id` in order; `SearchPage` adds the answer's `found`. A weight list whose length differs from `QueryBy`, or a `PerPage` above `typesense.MaxPerPage` (250, Typesense's limit), is an error before any request; callers page instead. - Ids and index names are path-escaped. - A non-2xx answer is a `typesense.StatusError` with the method, path and status, never the answer body. @@ -133,6 +135,10 @@ ids, err := svc.Engine().SearchIDs(ctx, svc.IndexName(&models.Post{}), beachcomb | `beachcomber.PageSearcher` | Optional engine interface: `SearchPage(ctx, index, q)` returns a page of ids and the engine's found count. | | `beachcomber.SearchResult` | `IDs` (candidates) and `Found` (documents the engine matched). | | `beachcomber.SearchPage(ctx, engine, index, q)` | Calls `PageSearcher` when the engine has it, else `SearchIDs` with `Found` set to the number of ids. | +| `beachcomber.IndexDropper` | Optional engine interface: `DropIndex(ctx, index)` drops the index and reports whether it existed. | +| `beachcomber.DropIndex(ctx, engine, index)` | Calls `IndexDropper` when the engine has it, else `Flush` and reports existed true. | +| `beachcomber.IndexEnsurer` | Optional engine interface: `EnsureIndex(ctx, index, schema)` creates a missing index, empty, from its schema. | +| `beachcomber.EnsureIndex(ctx, engine, index, schema)` | Calls `IndexEnsurer` when the engine has it; a no-op otherwise. | | `beachcomber.Gate`, `beachcomber.GateFunc` | The application kill-switch: `Enabled(ctx, db) bool`. | | `beachcomber.EngineFactory`, `beachcomber.RegisterEngine(name, factory)` | Registers an engine from an `init` function. | | `beachcomber.NullEngine`, `beachcomber.DefaultDriver` | The name of the built-in engine that indexes nothing, and the default of `search.driver`. | @@ -143,7 +149,7 @@ ids, err := svc.Engine().SearchIDs(ctx, svc.IndexName(&models.Post{}), beachcomb | Identifier | Description | |------------|-------------| | `typesense.Config`, `typesense.LoadConfig` | The `search.typesense.*` settings with their defaults; `BaseURL` is `{protocol}://{host}:{port}{path}`. | -| `typesense.Engine`, `typesense.New` | The `beachcomber.Engine` and `beachcomber.PageSearcher`, with `Config`. | +| `typesense.Engine`, `typesense.New` | The `beachcomber.Engine`, `beachcomber.PageSearcher`, `beachcomber.IndexDropper` and `beachcomber.IndexEnsurer`, with `Config`. | | `typesense.MaxPerPage` | 250, the largest page a search may ask for. | | `typesense.StatusError` | A non-2xx answer: `Method`, `Path`, `Code` and `StatusCode()`. | | `typesense.DriverName` | `typesense`. | diff --git a/modules/beachcomber/searchable.go b/modules/beachcomber/searchable.go index d9c48e2..4a14776 100644 --- a/modules/beachcomber/searchable.go +++ b/modules/beachcomber/searchable.go @@ -106,6 +106,43 @@ func SearchPage(ctx context.Context, e Engine, index string, q Query) (SearchRes return SearchResult{IDs: ids, Found: len(ids)}, nil } +// IndexDropper is implemented by an engine that can tell whether the index +// it dropped existed. It is optional; use DropIndex to call it. +type IndexDropper interface { + DropIndex(ctx context.Context, index string) (existed bool, err error) +} + +// DropIndex drops index on e and reports whether it existed: through +// IndexDropper when e implements it, otherwise through Flush, in which case +// it reports existed true because Flush cannot tell. +func DropIndex(ctx context.Context, e Engine, index string) (bool, error) { + if d, ok := e.(IndexDropper); ok { + return d.DropIndex(ctx, index) + } + if err := e.Flush(ctx, index); err != nil { + return false, err + } + return true, nil +} + +// IndexEnsurer is implemented by an engine that can create an empty index +// from its schema. It is optional; use EnsureIndex to call it. +type IndexEnsurer interface { + EnsureIndex(ctx context.Context, index string, schema map[string]any) error +} + +// EnsureIndex creates index on e from schema when it does not exist yet, +// through IndexEnsurer when e implements it; for other engines it does +// nothing. A full reindex calls it first so an empty source table still +// leaves an (empty) index behind, because Upsert of zero documents creates +// nothing. +func EnsureIndex(ctx context.Context, e Engine, index string, schema map[string]any) error { + if en, ok := e.(IndexEnsurer); ok { + return en.EnsureIndex(ctx, index, schema) + } + return nil +} + // Gate is the application kill-switch consulted before every sync. It // reads its setting through db, a clean session on the write's connection // (inside a caller's plain transaction, a savepoint of it); any error must diff --git a/modules/beachcomber/searchpage_test.go b/modules/beachcomber/searchpage_test.go index 9b85096..71fc435 100644 --- a/modules/beachcomber/searchpage_test.go +++ b/modules/beachcomber/searchpage_test.go @@ -78,3 +78,70 @@ func TestSearchPageNullEngine(t *testing.T) { t.Fatalf("res = %+v err = %v", res, err) } } + +// flushEngine records Flush calls and can fail them. +type flushEngine struct { + idsOnlyEngine + flushed []string + err error +} + +func (e *flushEngine) Flush(_ context.Context, index string) error { + e.flushed = append(e.flushed, index) + return e.err +} + +// dropperEngine implements IndexDropper and IndexEnsurer. +type dropperEngine struct { + flushEngine + existing map[string]bool + ensured []string +} + +func (e *dropperEngine) DropIndex(_ context.Context, index string) (bool, error) { + existed := e.existing[index] + delete(e.existing, index) + return existed, nil +} + +func (e *dropperEngine) EnsureIndex(_ context.Context, index string, _ map[string]any) error { + e.ensured = append(e.ensured, index) + e.existing[index] = true + return nil +} + +func TestDropIndex(t *testing.T) { + ctx := context.Background() + + fallback := &flushEngine{} + existed, err := DropIndex(ctx, fallback, "albums") + if err != nil || !existed || !reflect.DeepEqual(fallback.flushed, []string{"albums"}) { + t.Fatalf("fallback: existed=%v err=%v flushed=%v", existed, err, fallback.flushed) + } + boom := errors.New("boom") + failing := &flushEngine{err: boom} + if existed, err := DropIndex(ctx, failing, "albums"); !errors.Is(err, boom) || existed { + t.Fatalf("fallback error: existed=%v err=%v", existed, err) + } + + d := &dropperEngine{existing: map[string]bool{"albums": true}} + if existed, err := DropIndex(ctx, d, "albums"); err != nil || !existed { + t.Fatalf("dropper first: existed=%v err=%v", existed, err) + } + if existed, err := DropIndex(ctx, d, "albums"); err != nil || existed { + t.Fatalf("dropper second: existed=%v err=%v, want already absent", existed, err) + } + if len(d.flushed) != 0 { + t.Fatalf("dropper path must not call Flush: %v", d.flushed) + } + + if err := EnsureIndex(ctx, d, "albums", map[string]any{"fields": []any{}}); err != nil || !reflect.DeepEqual(d.ensured, []string{"albums"}) { + t.Fatalf("EnsureIndex: err=%v ensured=%v", err, d.ensured) + } + if err := EnsureIndex(ctx, fallback, "albums", nil); err != nil { + t.Fatalf("EnsureIndex without IndexEnsurer must be a no-op: %v", err) + } + if err := EnsureIndex(ctx, nullEngine{}, "albums", nil); err != nil { + t.Fatalf("null engine: %v", err) + } +} diff --git a/modules/beachcomber/typesense/engine.go b/modules/beachcomber/typesense/engine.go index dfd428e..d7e8a9d 100644 --- a/modules/beachcomber/typesense/engine.go +++ b/modules/beachcomber/typesense/engine.go @@ -47,6 +47,10 @@ type Engine struct { } var _ beachcomber.Engine = (*Engine)(nil) +var ( + _ beachcomber.IndexDropper = (*Engine)(nil) + _ beachcomber.IndexEnsurer = (*Engine)(nil) +) // New returns an engine for cfg. Every request is bounded by // cfg.ConnectionTimeout. @@ -214,6 +218,33 @@ func (e *Engine) Flush(ctx context.Context, index string) error { return nil } +// DropIndex drops the collection with DELETE /collections/{index} and +// reports whether it existed: 2xx is true, 404 is false, any other status is +// a StatusError. It implements beachcomber.IndexDropper. +func (e *Engine) DropIndex(ctx context.Context, index string) (bool, error) { + path := "/collections/" + url.PathEscape(index) + code, _, err := e.do(ctx, http.MethodDelete, path, nil, nil, "") + if err != nil { + return false, err + } + switch { + case ok2xx(code): + return true, nil + case code == http.StatusNotFound: + return false, nil + default: + return false, statusError(http.MethodDelete, path, code) + } +} + +// EnsureIndex creates the collection from schema when Typesense does not +// have it (GET /collections/{index}, then POST /collections on 404, 409 as +// success), the same step Upsert runs. It implements +// beachcomber.IndexEnsurer. +func (e *Engine) EnsureIndex(ctx context.Context, index string, schema map[string]any) error { + return e.ensureCollection(ctx, index, schema) +} + // searchAnswer is the part of a search answer SearchIDs and SearchPage read. type searchAnswer struct { Found int `json:"found"` diff --git a/modules/beachcomber/typesense/engine_test.go b/modules/beachcomber/typesense/engine_test.go index 23457fe..077c5c3 100644 --- a/modules/beachcomber/typesense/engine_test.go +++ b/modules/beachcomber/typesense/engine_test.go @@ -367,3 +367,66 @@ func TestSyncEngineRegistration(t *testing.T) { t.Fatal("a blank API key counts as configured") } } + +func TestEngineDropIndex(t *testing.T) { + f, e := newFake(t) + f.answer("DELETE /ts/collections/albums", http.StatusOK, `{"name":"albums"}`) + existed, err := e.DropIndex(t.Context(), "albums") + if err != nil || !existed { + t.Fatalf("200: existed=%v err=%v", existed, err) + } + f.answer("DELETE /ts/collections/albums", http.StatusNotFound, `{"message":"Not Found"}`) + existed, err = e.DropIndex(t.Context(), "albums") + if err != nil || existed { + t.Fatalf("404: existed=%v err=%v, want already absent", existed, err) + } + f.answer("DELETE /ts/collections/albums", http.StatusInternalServerError, `{"message":"boom"}`) + if _, err := e.DropIndex(t.Context(), "albums"); statusOf(err) != http.StatusInternalServerError { + t.Fatalf("500: err=%v", err) + } + calls := f.take() + if len(calls) != 3 { + t.Fatalf("calls = %+v", calls) + } + for _, c := range calls { + if c.method != http.MethodDelete || c.path != "/ts/collections/albums" || c.key != testKey { + t.Fatalf("call = %+v", c) + } + } + existed, err = beachcomber.DropIndex(t.Context(), e, "a b") + if err != nil || existed { + t.Fatalf("through beachcomber.DropIndex: existed=%v err=%v", existed, err) + } + if got := f.take(); len(got) != 1 || got[0].path != "/ts/collections/a%20b" { + t.Fatalf("escaped path calls = %+v", got) + } +} + +func TestEngineEnsureIndex(t *testing.T) { + f, e := newFake(t) + f.answer("POST /ts/collections", http.StatusCreated, `{"name":"albums"}`) + schema := map[string]any{"fields": []map[string]any{{"name": "title", "type": "string"}}} + if err := beachcomber.EnsureIndex(t.Context(), e, "albums", schema); err != nil { + t.Fatal(err) + } + calls := f.take() + if len(calls) != 2 || calls[0].method != http.MethodGet || calls[1].method != http.MethodPost || calls[1].path != "/ts/collections" { + t.Fatalf("first ensure calls = %+v", calls) + } + var created map[string]any + if err := json.Unmarshal([]byte(calls[1].body), &created); err != nil { + t.Fatal(err) + } + if created["name"] != "albums" || created["fields"] == nil { + t.Fatalf("create body = %s", calls[1].body) + } + + f.answer("GET /ts/collections/albums", http.StatusOK, `{"name":"albums"}`) + if err := e.EnsureIndex(t.Context(), "albums", schema); err != nil { + t.Fatal(err) + } + calls = f.take() + if len(calls) != 1 || calls[0].method != http.MethodGet { + t.Fatalf("second ensure must send no create: %+v", calls) + } +}