diff --git a/modules/beachcomber/README.md b/modules/beachcomber/README.md index 6d4df4f..be9d2dd 100644 --- a/modules/beachcomber/README.md +++ b/modules/beachcomber/README.md @@ -21,7 +21,27 @@ The package itself knows no search server. An engine package registers itself fr - The `beachcomber.Engine` driver contract: `Name`, `Configured`, `Upsert`, `Delete`, `Flush` and `SearchIDs` with a `beachcomber.Query`. - 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`): the `X-TYPESENSE-API-KEY` header on every request; `Upsert` reads the collection and creates it from the schema on 404, then imports JSON lines with `action=upsert`; `Delete` and `Flush` treat 404 as success; `SearchIDs` returns `hits[].document.id`. +- 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. + - `SearchIDs` sends `q` (default `*`), `query_by`, `filter_by`, `sort_by`, `page` and `per_page`, and returns `hits[].document.id` in order. + - 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. + +## Sync semantics + +- **After commit.** The callbacks register the sync with `lagoon.AfterCommit`. Inside `lagoon.Transaction` it runs after that transaction commits, and not at all when it rolls back. A single-statement write, for which GORM opens its own transaction, syncs after that commit and not when the write fails. Inside a plain `gorm` transaction there is no commit hook, so the sync runs immediately through the transaction's handle. Its reads run in a savepoint, so a failed read never aborts the caller's transaction. +- **Inline and non-fatal.** The sync runs in the writing goroutine, after the commit, so a create followed by a search sees the document. Every engine request is bounded by the engine's timeout (`search.typesense.connection_timeout_seconds`), and the caller's context cancellation does not abandon it. A failure, a timeout or a panic is logged at Warn as `search: sync failed` with the index, key and operation. The write is already committed and stays so. The log never carries the document or the API key. +- **Three gates, before any request.** Nothing is sent when: + 1. the engine is not configured (the `null` engine, or Typesense with an empty `search.typesense.api_key`); + 2. no `*gorm.DB` is published on the app (a fresh install); + 3. the application `beachcomber.Gate` reports off. A gate must treat a read error as off. + + When the engine is not configured, the callbacks do not even register work. +- **Reload, then decide.** The row is reloaded by primary key, including soft-deleted rows. A delete, a row that is gone, a soft-deleted row (a set `gorm.DeletedAt`) or a row whose `ShouldBeSearchable` is false has its document deleted. Restoring a soft-deleted row is an ordinary update and indexes it again. An error from `ToSearchableArray` is logged and nothing is sent, which lets a model refuse a document that would break scoping. A document without an `id` gets the key. +- **Rows only.** A statement without a primary key value, such as `Model(&T{}).Where(…).Updates(…)` or `Delete(&T{}, id)`, cannot be synced row by row and is skipped. Bulk paths call `beachcomber.Service.Sync` or `beachcomber.Service.Remove` per row, or reindex. +- **Candidates, not answers.** `beachcomber.Engine.SearchIDs` returns candidate ids from an external index that may be stale. Callers must re-gate every id in SQL (ownership, visibility, soft deletes) before they expose a row. An empty result is an empty list, never an error. ## Usage @@ -120,6 +140,7 @@ ids, err := svc.Engine().SearchIDs(ctx, svc.IndexName(&models.Post{}), beachcomb |------------|-------------| | `typesense.Config`, `typesense.LoadConfig` | The `search.typesense.*` settings with their defaults; `BaseURL` is `{protocol}://{host}:{port}{path}`. | | `typesense.Engine`, `typesense.New` | The `beachcomber.Engine`, with `Config`. | +| `typesense.StatusError` | A non-2xx answer: `Method`, `Path`, `Code` and `StatusCode()`. | | `typesense.DriverName` | `typesense`. | | `typesense.DefaultHost`, `typesense.DefaultPort`, `typesense.DefaultProtocol`, `typesense.DefaultConnectionTimeout`, `typesense.DefaultImportAction` | Defaults of the configuration keys. | diff --git a/modules/beachcomber/beachcomber.go b/modules/beachcomber/beachcomber.go index 5138f40..800bcbf 100644 --- a/modules/beachcomber/beachcomber.go +++ b/modules/beachcomber/beachcomber.go @@ -138,27 +138,33 @@ func (s *Service) Logger() *slog.Logger { return s.log } +// commitCallback is GORM's commit of a transaction it opened itself. +const commitCallback = "gorm:commit_or_rollback_transaction" + // installCallbacks registers the sync callbacks on gdb, replacing earlier // ones so a handle shared by several apps syncs through the most recent -// service. +// service. Each runs after the model's own after hook and before GORM +// commits a single-statement write: GORM appends a callback that names only +// an After anchor to the end of the chain, past the commit and past +// lagoon's after-commit flush, where the registered work would never run. func (s *Service) installCallbacks(gdb *gorm.DB) error { cb := gdb.Callback() if cb.Create().Get(CallbackAfterCreate) == nil { - if err := cb.Create().After("gorm:after_create").Register(CallbackAfterCreate, s.afterCreate); err != nil { + if err := cb.Create().After("gorm:after_create").Before(commitCallback).Register(CallbackAfterCreate, s.afterCreate); err != nil { return err } } else if err := cb.Create().Replace(CallbackAfterCreate, s.afterCreate); err != nil { return err } if cb.Update().Get(CallbackAfterUpdate) == nil { - if err := cb.Update().After("gorm:after_update").Register(CallbackAfterUpdate, s.afterUpdate); err != nil { + if err := cb.Update().After("gorm:after_update").Before(commitCallback).Register(CallbackAfterUpdate, s.afterUpdate); err != nil { return err } } else if err := cb.Update().Replace(CallbackAfterUpdate, s.afterUpdate); err != nil { return err } if cb.Delete().Get(CallbackAfterDelete) == nil { - if err := cb.Delete().After("gorm:after_delete").Register(CallbackAfterDelete, s.afterDelete); err != nil { + if err := cb.Delete().After("gorm:after_delete").Before(commitCallback).Register(CallbackAfterDelete, s.afterDelete); err != nil { return err } } else if err := cb.Delete().Replace(CallbackAfterDelete, s.afterDelete); err != nil { diff --git a/modules/beachcomber/searchable.go b/modules/beachcomber/searchable.go index b49877f..3b5549d 100644 --- a/modules/beachcomber/searchable.go +++ b/modules/beachcomber/searchable.go @@ -74,7 +74,9 @@ type Query struct { } // Gate is the application kill-switch consulted before every sync. It -// reads its setting through db; any error must count as off. +// 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 +// count as off. type Gate interface { Enabled(ctx context.Context, db *gorm.DB) bool } diff --git a/modules/beachcomber/sync.go b/modules/beachcomber/sync.go index b263deb..8f57423 100644 --- a/modules/beachcomber/sync.go +++ b/modules/beachcomber/sync.go @@ -182,7 +182,7 @@ func (s *Service) syncOne(ctx context.Context, db *gorm.DB, p pending) (err erro return nil } ctx = context.WithoutCancel(ctx) - sess := db.Session(&gorm.Session{NewDB: true, Context: ctx}) + sess := cleanSession(db, ctx) var ( skip bool @@ -249,6 +249,16 @@ func (s *Service) syncOne(ctx context.Context, db *gorm.DB, p pending) (err erro return nil } +// cleanSession returns a handle on db's connection (the pool, or the open +// transaction) with an empty statement. db.Session with NewDB and a Context +// is not enough on a callback's handle: the Context makes it clone the +// write's statement (model, table, clauses), and a later WithContext on the +// result continues from that clone, so a gate query would run against the +// written model's table. +func cleanSession(db *gorm.DB, ctx context.Context) *gorm.DB { + return db.Session(&gorm.Session{NewDB: true, Context: ctx}).Clauses().Session(&gorm.Session{NewDB: true}) +} + // inSavepoint runs fn on db. Inside a transaction fn runs in a savepoint, // so a failed read (a table that does not exist yet, say) is rolled back // to it and never aborts the caller's transaction. diff --git a/modules/beachcomber/typesense/engine.go b/modules/beachcomber/typesense/engine.go index 4ed2a9b..4e53e69 100644 --- a/modules/beachcomber/typesense/engine.go +++ b/modules/beachcomber/typesense/engine.go @@ -317,6 +317,21 @@ func unwrapURLError(err error) error { func ok2xx(code int) bool { return code >= 200 && code < 300 } -func statusError(method, path string, code int) error { - return fmt.Errorf("typesense: %s %s: status %d", method, path, code) +// StatusError is a Typesense answer outside 2xx. It carries the method, +// path and status, never the answer body, which could echo request data. +type StatusError struct { + Method string + Path string + Code int +} + +func (e *StatusError) Error() string { + return fmt.Sprintf("typesense: %s %s: status %d", e.Method, e.Path, e.Code) +} + +// StatusCode returns the HTTP status Typesense answered with. +func (e *StatusError) StatusCode() int { return e.Code } + +func statusError(method, path string, code int) error { + return &StatusError{Method: method, Path: path, Code: code} }