- 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
359 lines
9.5 KiB
Go
359 lines
9.5 KiB
Go
package beachcomber
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"reflect"
|
|
"strconv"
|
|
|
|
"git.golem15.com/golem15/summercms/modules/lagoon"
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/schema"
|
|
)
|
|
|
|
// operation is what a sync does with the written row's document.
|
|
type operation string
|
|
|
|
const (
|
|
opUpsert operation = "upsert"
|
|
opDelete operation = "delete"
|
|
)
|
|
|
|
const syncSavepoint = "beachcomber_sync"
|
|
|
|
var (
|
|
searchableType = reflect.TypeFor[Searchable]()
|
|
deletedAtType = reflect.TypeFor[gorm.DeletedAt]()
|
|
)
|
|
|
|
// pending is one written row waiting for its after-commit sync. It holds
|
|
// the primary key, never the model, so the sync always reads the
|
|
// committed row.
|
|
type pending struct {
|
|
sch *schema.Schema
|
|
typ reflect.Type
|
|
pk any
|
|
key string
|
|
index string
|
|
op operation
|
|
}
|
|
|
|
func (s *Service) afterCreate(db *gorm.DB) { s.afterWrite(db, opUpsert) }
|
|
func (s *Service) afterUpdate(db *gorm.DB) { s.afterWrite(db, opUpsert) }
|
|
func (s *Service) afterDelete(db *gorm.DB) { s.afterWrite(db, opDelete) }
|
|
|
|
// afterWrite registers one after-commit sync per Searchable row of the
|
|
// statement. A statement without a primary key value (a batch update or
|
|
// delete through an empty model) is skipped: no document can be built for
|
|
// it. Nothing is registered while the engine is not configured.
|
|
func (s *Service) afterWrite(db *gorm.DB, op operation) {
|
|
if db.Error != nil || db.Statement == nil || db.Statement.Schema == nil {
|
|
return
|
|
}
|
|
if s.engine == nil || !s.engine.Configured() {
|
|
return
|
|
}
|
|
sch := db.Statement.Schema
|
|
if !reflect.PointerTo(sch.ModelType).Implements(searchableType) {
|
|
return
|
|
}
|
|
pkField := sch.PrioritizedPrimaryField
|
|
if pkField == nil {
|
|
return
|
|
}
|
|
ctx := db.Statement.Context
|
|
if ctx == nil {
|
|
ctx = context.Background()
|
|
}
|
|
index := s.IndexName(reflect.New(sch.ModelType).Interface().(Searchable))
|
|
add := func(v reflect.Value) {
|
|
for v.Kind() == reflect.Pointer {
|
|
if v.IsNil() {
|
|
return
|
|
}
|
|
v = v.Elem()
|
|
}
|
|
if v.Kind() != reflect.Struct || v.Type() != sch.ModelType {
|
|
return
|
|
}
|
|
pk, zero := pkField.ValueOf(ctx, v)
|
|
if zero {
|
|
return
|
|
}
|
|
p := pending{sch: sch, typ: sch.ModelType, pk: pk, key: searchKey(v, pk), index: index, op: op}
|
|
lagoon.AfterCommit(ctx, db, func(ctx context.Context, db *gorm.DB) {
|
|
if err := s.syncOne(ctx, db, p); err != nil {
|
|
s.warn(p, err)
|
|
}
|
|
})
|
|
}
|
|
rv := db.Statement.ReflectValue
|
|
switch rv.Kind() {
|
|
case reflect.Slice, reflect.Array:
|
|
for i := 0; i < rv.Len(); i++ {
|
|
add(rv.Index(i))
|
|
}
|
|
default:
|
|
add(rv)
|
|
}
|
|
}
|
|
|
|
// Sync upserts the document of model (a pointer to a Searchable struct
|
|
// with its primary key set) through the same path as the callbacks: the
|
|
// gates, a reload by primary key, and a delete when the row is gone, soft
|
|
// deleted or not searchable. db may be nil to use the published handle. A
|
|
// skipped sync returns nil.
|
|
func (s *Service) Sync(ctx context.Context, db *gorm.DB, model any) error {
|
|
return s.explicit(ctx, db, model, opUpsert)
|
|
}
|
|
|
|
// Remove deletes the document of model (a pointer to a Searchable struct
|
|
// with its primary key set) through the same gated path as Sync.
|
|
func (s *Service) Remove(ctx context.Context, db *gorm.DB, model any) error {
|
|
return s.explicit(ctx, db, model, opDelete)
|
|
}
|
|
|
|
func (s *Service) explicit(ctx context.Context, db *gorm.DB, model any, op operation) error {
|
|
if s == nil {
|
|
return fmt.Errorf("beachcomber: nil service")
|
|
}
|
|
if ctx == nil {
|
|
ctx = context.Background()
|
|
}
|
|
if db == nil {
|
|
published, ok := s.publishedDB()
|
|
if !ok {
|
|
return nil
|
|
}
|
|
db = published
|
|
}
|
|
rv := reflect.ValueOf(model)
|
|
if rv.Kind() != reflect.Pointer || rv.IsNil() || rv.Elem().Kind() != reflect.Struct {
|
|
return fmt.Errorf("beachcomber: %T is not a pointer to a model struct", model)
|
|
}
|
|
if _, ok := model.(Searchable); !ok {
|
|
return fmt.Errorf("beachcomber: %T does not implement Searchable", model)
|
|
}
|
|
stmt := &gorm.Statement{DB: db}
|
|
if err := stmt.Parse(model); err != nil {
|
|
return fmt.Errorf("beachcomber: parse %T: %w", model, err)
|
|
}
|
|
pkField := stmt.Schema.PrioritizedPrimaryField
|
|
if pkField == nil {
|
|
return fmt.Errorf("beachcomber: %T has no primary key", model)
|
|
}
|
|
v := rv.Elem()
|
|
pk, zero := pkField.ValueOf(ctx, v)
|
|
if zero {
|
|
return fmt.Errorf("beachcomber: %T has a zero primary key", model)
|
|
}
|
|
p := pending{
|
|
sch: stmt.Schema,
|
|
typ: stmt.Schema.ModelType,
|
|
pk: pk,
|
|
key: searchKey(v, pk),
|
|
index: s.IndexName(model.(Searchable)),
|
|
op: op,
|
|
}
|
|
return s.syncOne(ctx, db, p)
|
|
}
|
|
|
|
// syncOne runs the gates, in order and without a request: the engine is
|
|
// configured, a database is published, the application Gate is on. It then
|
|
// reloads the row by primary key and upserts its document, or deletes it
|
|
// when the operation is a delete or the row is gone, soft deleted or not
|
|
// searchable. The engine call is bounded by the engine's own timeout; the
|
|
// caller's cancellation does not abandon it.
|
|
func (s *Service) syncOne(ctx context.Context, db *gorm.DB, p pending) (err error) {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
err = fmt.Errorf("panic: %v", r)
|
|
}
|
|
}()
|
|
if s.engine == nil || !s.engine.Configured() {
|
|
return nil
|
|
}
|
|
if _, ok := s.publishedDB(); !ok {
|
|
return nil
|
|
}
|
|
if db == nil {
|
|
return nil
|
|
}
|
|
ctx = context.WithoutCancel(ctx)
|
|
sess := cleanSession(db, ctx)
|
|
|
|
var (
|
|
skip bool
|
|
remove bool
|
|
doc map[string]any
|
|
idx map[string]any
|
|
)
|
|
err = inSavepoint(sess, func(tx *gorm.DB) error {
|
|
if g := s.currentGate(); g != nil && !g.Enabled(ctx, tx) {
|
|
skip = true
|
|
return nil
|
|
}
|
|
if p.op == opDelete {
|
|
remove = true
|
|
return nil
|
|
}
|
|
model := reflect.New(p.typ)
|
|
pkCol := tx.Statement.Quote(p.sch.PrioritizedPrimaryField.DBName)
|
|
err := tx.Unscoped().Where(pkCol+" = ?", p.pk).Take(model.Interface()).Error
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
remove = true
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return fmt.Errorf("reload: %w", err)
|
|
}
|
|
if softDeleted(ctx, p.sch, model.Elem()) {
|
|
remove = true
|
|
return nil
|
|
}
|
|
m := model.Interface().(Searchable)
|
|
if !m.ShouldBeSearchable() {
|
|
remove = true
|
|
return nil
|
|
}
|
|
doc, err = m.ToSearchableArray(ctx, tx)
|
|
if err != nil {
|
|
return fmt.Errorf("build document: %w", err)
|
|
}
|
|
if len(doc) == 0 {
|
|
skip = true
|
|
return nil
|
|
}
|
|
if _, ok := doc["id"]; !ok {
|
|
doc["id"] = p.key
|
|
}
|
|
if sp, ok := model.Interface().(IndexSchemaProvider); ok {
|
|
idx = sp.SearchIndexSchema()
|
|
}
|
|
return nil
|
|
})
|
|
if err != nil || skip {
|
|
return err
|
|
}
|
|
if remove {
|
|
if err := s.engine.Delete(ctx, p.index, []string{p.key}); err != nil {
|
|
return fmt.Errorf("delete: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
if err := s.engine.Upsert(ctx, p.index, idx, []map[string]any{doc}); err != nil {
|
|
return fmt.Errorf("upsert: %w", err)
|
|
}
|
|
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.
|
|
func inSavepoint(db *gorm.DB, fn func(tx *gorm.DB) error) error {
|
|
if _, inTx := db.Statement.ConnPool.(gorm.TxCommitter); !inTx {
|
|
return fn(db)
|
|
}
|
|
if err := db.SavePoint(syncSavepoint).Error; err != nil {
|
|
return fmt.Errorf("savepoint: %w", err)
|
|
}
|
|
var err error
|
|
func() {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
err = fmt.Errorf("panic: %v", r)
|
|
}
|
|
}()
|
|
err = fn(db)
|
|
}()
|
|
if err != nil {
|
|
db.RollbackTo(syncSavepoint)
|
|
return err
|
|
}
|
|
db.Exec("RELEASE SAVEPOINT " + syncSavepoint)
|
|
return nil
|
|
}
|
|
|
|
// softDeleted reports whether any gorm.DeletedAt field of v is set.
|
|
func softDeleted(ctx context.Context, sch *schema.Schema, v reflect.Value) bool {
|
|
for _, f := range sch.Fields {
|
|
if f.FieldType != deletedAtType {
|
|
continue
|
|
}
|
|
val, zero := f.ValueOf(ctx, v)
|
|
if zero {
|
|
continue
|
|
}
|
|
if d, ok := val.(gorm.DeletedAt); ok && d.Valid {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// searchKey is SearchKeyer.SearchKey when the model implements it and
|
|
// returns a key, else the decimal primary key.
|
|
func searchKey(v reflect.Value, pk any) string {
|
|
if v.CanAddr() {
|
|
if k, ok := v.Addr().Interface().(SearchKeyer); ok {
|
|
if key := k.SearchKey(); key != "" {
|
|
return key
|
|
}
|
|
}
|
|
} else if k, ok := v.Interface().(SearchKeyer); ok {
|
|
if key := k.SearchKey(); key != "" {
|
|
return key
|
|
}
|
|
}
|
|
switch n := pk.(type) {
|
|
case uint:
|
|
return strconv.FormatUint(uint64(n), 10)
|
|
case uint32:
|
|
return strconv.FormatUint(uint64(n), 10)
|
|
case uint64:
|
|
return strconv.FormatUint(n, 10)
|
|
case int:
|
|
return strconv.Itoa(n)
|
|
case int32:
|
|
return strconv.FormatInt(int64(n), 10)
|
|
case int64:
|
|
return strconv.FormatInt(n, 10)
|
|
default:
|
|
return fmt.Sprint(pk)
|
|
}
|
|
}
|
|
|
|
func (s *Service) publishedDB() (*gorm.DB, bool) {
|
|
if s.app == nil {
|
|
return nil, false
|
|
}
|
|
gdb, ok := s.app.Lookup[*gorm.DB]()
|
|
if !ok || gdb == nil {
|
|
return nil, false
|
|
}
|
|
return gdb, true
|
|
}
|
|
|
|
// warn logs a failed sync with the index, key and operation. It never logs
|
|
// the document or the engine's credentials.
|
|
func (s *Service) warn(p pending, err error) {
|
|
s.Logger().Warn("search: sync failed",
|
|
slog.String("index", p.index),
|
|
slog.String("key", p.key),
|
|
slog.String("operation", string(p.op)),
|
|
slog.String("error", err.Error()),
|
|
)
|
|
}
|