fix(11-05): sync single-statement and plain-transaction writes, type engine status errors
- pin the sync callbacks before gorm:commit_or_rollback_transaction: an After-only anchor is appended past the commit and lagoon's after-commit flush, so single-statement writes never synced - give the gate and the document builder a clean session: Session with NewDB and a Context clones the write's statement, and a later WithContext queried through the written model's table - typesense.StatusError carries method, path and status, never the body - README: sync semantics, the three gates, delete on soft delete, and the SQL re-gate required of SearchIDs callers
This commit is contained in:
@@ -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`.
|
- 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`.
|
- 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.
|
- `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
|
## 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.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.Engine`, `typesense.New` | The `beachcomber.Engine`, with `Config`. |
|
||||||
|
| `typesense.StatusError` | A non-2xx answer: `Method`, `Path`, `Code` and `StatusCode()`. |
|
||||||
| `typesense.DriverName` | `typesense`. |
|
| `typesense.DriverName` | `typesense`. |
|
||||||
| `typesense.DefaultHost`, `typesense.DefaultPort`, `typesense.DefaultProtocol`, `typesense.DefaultConnectionTimeout`, `typesense.DefaultImportAction` | Defaults of the configuration keys. |
|
| `typesense.DefaultHost`, `typesense.DefaultPort`, `typesense.DefaultProtocol`, `typesense.DefaultConnectionTimeout`, `typesense.DefaultImportAction` | Defaults of the configuration keys. |
|
||||||
|
|
||||||
|
|||||||
@@ -138,27 +138,33 @@ func (s *Service) Logger() *slog.Logger {
|
|||||||
return s.log
|
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
|
// installCallbacks registers the sync callbacks on gdb, replacing earlier
|
||||||
// ones so a handle shared by several apps syncs through the most recent
|
// 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 {
|
func (s *Service) installCallbacks(gdb *gorm.DB) error {
|
||||||
cb := gdb.Callback()
|
cb := gdb.Callback()
|
||||||
if cb.Create().Get(CallbackAfterCreate) == nil {
|
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
|
return err
|
||||||
}
|
}
|
||||||
} else if err := cb.Create().Replace(CallbackAfterCreate, s.afterCreate); err != nil {
|
} else if err := cb.Create().Replace(CallbackAfterCreate, s.afterCreate); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if cb.Update().Get(CallbackAfterUpdate) == nil {
|
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
|
return err
|
||||||
}
|
}
|
||||||
} else if err := cb.Update().Replace(CallbackAfterUpdate, s.afterUpdate); err != nil {
|
} else if err := cb.Update().Replace(CallbackAfterUpdate, s.afterUpdate); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if cb.Delete().Get(CallbackAfterDelete) == nil {
|
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
|
return err
|
||||||
}
|
}
|
||||||
} else if err := cb.Delete().Replace(CallbackAfterDelete, s.afterDelete); err != nil {
|
} else if err := cb.Delete().Replace(CallbackAfterDelete, s.afterDelete); err != nil {
|
||||||
|
|||||||
@@ -74,7 +74,9 @@ type Query struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Gate is the application kill-switch consulted before every sync. It
|
// 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 {
|
type Gate interface {
|
||||||
Enabled(ctx context.Context, db *gorm.DB) bool
|
Enabled(ctx context.Context, db *gorm.DB) bool
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -182,7 +182,7 @@ func (s *Service) syncOne(ctx context.Context, db *gorm.DB, p pending) (err erro
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
ctx = context.WithoutCancel(ctx)
|
ctx = context.WithoutCancel(ctx)
|
||||||
sess := db.Session(&gorm.Session{NewDB: true, Context: ctx})
|
sess := cleanSession(db, ctx)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
skip bool
|
skip bool
|
||||||
@@ -249,6 +249,16 @@ func (s *Service) syncOne(ctx context.Context, db *gorm.DB, p pending) (err erro
|
|||||||
return nil
|
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,
|
// 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
|
// so a failed read (a table that does not exist yet, say) is rolled back
|
||||||
// to it and never aborts the caller's transaction.
|
// to it and never aborts the caller's transaction.
|
||||||
|
|||||||
@@ -317,6 +317,21 @@ func unwrapURLError(err error) error {
|
|||||||
|
|
||||||
func ok2xx(code int) bool { return code >= 200 && code < 300 }
|
func ok2xx(code int) bool { return code >= 200 && code < 300 }
|
||||||
|
|
||||||
func statusError(method, path string, code int) error {
|
// StatusError is a Typesense answer outside 2xx. It carries the method,
|
||||||
return fmt.Errorf("typesense: %s %s: status %d", method, path, code)
|
// 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}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user