feat(14-01): beachcomber DropIndex and EnsureIndex for reindex tooling
- optional IndexDropper reports whether a dropped index existed; DropIndex falls back to Flush - optional IndexEnsurer creates an empty index from its schema; EnsureIndex is a no-op otherwise - typesense implements both (404 is already absent; ensure reuses collection creation)
This commit is contained in:
@@ -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`. |
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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"`
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user